RocketMQ 架构与术语详解
环境: KUBECONFIG=/bpx/.145-admin.conf
实例名称: rocketmq-ddffdb1a
命名空间: qfusion-admin
版本: RocketMQ 4.9.7
集群模式: DLedger (Raft 一致性协议)

一、部署架构总览
1.1 集群拓扑图
┌──────────────────────────────────────────────────────────────────────────────────┐
│ qfusion-admin Namespace │
│ (RocketMQ 实例: rocketmq-ddffdb1a) │
├──────────────────────────────────────────────────────────────────────────────────┤
│ │
│ ┌─────────────────────────────────────────────────────────────────────────┐ │
│ │ NameServer 集群 (3副本) │ ���
│ ├───────────────────┬───────────────────┬──────────────────────────────────┤ │
│ │ nameserver-0 │ nameserver-1 │ nameserver-2 │ │
│ │ qfusion4 │ qfusion2 │ qfusion1 │ │
│ │ 245.0.3.219:9876 │ 245.0.1.220:9876 │ 245.0.0.45:9876 │ │
│ └───────────────────┴───────────────────┴──────────────────────────────────┘ │
│ │ │
│ ▼ │
│ ┌─────────────────────────────────────────────────────────────────────────┐ │
│ │ DLedger Raft Group (rocketmq-ddffdb1a-0) │ │
│ │ (3副本高可用集群 - 一致性协议: Raft) │ │
│ ├───────────────────┬───────────────────┬──────────────────────────────────┤ │
│ │ Broker-0 (LEADER) │ Broker-1 (FOLLOWER)│ Broker-2 (FOLLOWER) │ │
│ │ qfusion4 │ qfusion1 │ qfusion2 │ │
│ │ 245.0.3.63 │ 245.0.0.157 │ 245.0.1.254 │ │
│ │ BrokerId=0 │ BrokerId=1 │ BrokerId=3 │ │
│ │ x.x.x.146:3 │ x.x.x.146:9 │ x.x.x.146:6 │ │
│ └───────────────────┴───────────────────┴──────────────────────────────────┘ │
│ │
│ ┌─────────────────────────────────────────────────────────────────────────┐ │
│ │ qfusion Namespace │ │
│ │ (管理平面组件) │ │
│ ├───────────────────┬──────────────────────────────────────────────────────┤ │
│ │ rocketmq-operator │ rocketmq-webserver (2副本) │ │
│ │ Pod 管理 │ Web 控制台 (9090端口) │ │
│ └───────────────────┴──────────────────────────────────────────────────────┘ │
│ │
└──────────────────────────────────────────────────────────────────────────────────┘
1.2 资源清单
| 资源类型 |
名称 |
命名空间 |
副本数 |
状态 |
| Broker |
rocketmq-ddffdb1a |
qfusion-admin |
3 (DLedger) |
Running |
| NameServer |
rocketmq-ddffdb1a-nameserver |
qfusion-admin |
3 |
Running |
| Operator |
rocketmq-operator |
qfusion |
1 |
Running |
| WebServer |
rocketmq-webserver |
qfusion |
2 |
Running |
1.3 组件分布
| 组件 |
Pod 名称 |
节点 |
Pod IP |
服务 IP |
| NameServer-0 |
rocketmq-ddffdb1a-nameserver-0-0 |
qfusion4 |
245.0.3.219 |
246.108.185.135 (Client) |
| NameServer-1 |
rocketmq-ddffdb1a-nameserver-1-0 |
qfusion2 |
245.0.1.220 |
246.108.185.135 (Client) |
| NameServer-2 |
rocketmq-ddffdb1a-nameserver-2-0 |
qfusion1 |
245.0.0.45 |
246.108.185.135 (Client) |
| Broker-0 (Leader) |
rocketmq-ddffdb1a-0-0-0 |
qfusion4 |
245.0.3.63 |
246.106.184.88 |
| Broker-1 (Follower) |
rocketmq-ddffdb1a-0-1-0 |
qfusion1 |
245.0.0.157 |
246.96.207.183 |
| Broker-2 (Follower) |
rocketmq-ddffdb1a-0-2-0 |
qfusion2 |
245.0.1.254 |
246.98.35.243 |
二、核心术语详解
2.1 架构组件术语
| 术语 |
英文 |
说明 |
| NameServer |
NameServer |
路由中心,维护 Topic/Broker 路由信息,无状态 |
| Broker |
Broker |
消息存储节点,负责收发、存储消息 |
| DLedger |
DLedger |
基于 Raft 协议的副本一致性机制 |
| Leader |
Leader |
DLedger 组的主节点,负责写入 |
| Follower |
Follower |
DLedger 组的从节点,同步复制数据 |
2.2 消息模型术语
| 术语 |
英文 |
说明 |
| Topic |
Topic |
消息主题,生产者发送消息的逻辑分类 |
| Queue |
MessageQueue |
Topic 的分区,一个 Topic 可有多个 Queue |
| BrokerId |
Broker ID |
Broker 在 Raft 组中的唯一标识 |
| Broker Name |
Broker Name |
Broker 组的名称,同一组内共享 |
2.3 消费者模型术语
| 术语 |
英文 |
说明 |
| ConsumerGroup |
Consumer Group |
消费者组,多个消费者组成逻辑分组 |
| Offset |
Consumer Offset |
消费进度,记录消费者消费到哪条消息 |
| Accumulation |
Message Lag |
消息堆积量 = BrokerOffset - ConsumerOffset |
| Rebalance |
Rebalance |
队列重新分配,消费者上下线时触发 |
2.4 性能指标术语
| 术语 |
英文 |
说明 |
| TPS |
Transactions Per Second |
每秒事务数(消息吞吐量) |
| InTPS |
Inbound TPS |
生产速率(每秒写入消息数) |
| OutTPS |
Outbound TPS |
消费速率(每秒消费消息数) |
| MaxOffset |
Max Offset |
Broker 当前最大消息偏移量 |
三、服务端口说明
3.1 NameServer 端口
| 端口 |
协议 |
说明 |
| 9876 |
TCP |
NameServer 服务端口,客户端连接端口 |
3.2 Broker 端口
| 端口 |
协议 |
说明 |
| 10911 |
TCP |
Broker 主服务端口,消息收发 |
| 10909 |
TCP |
Broker VIP 端口,内部通信 |
| 10912 |
TCP |
Broker HA 端口,高可用通信 |
| 40911 |
TCP |
DLedger Raft 协议端口,副本间通信 |
3.3 管理服务端口
| 端口 |
协议 |
说明 |
| 9090 |
HTTP |
RocketMQ Web 控制台 |
| 5557 |
HTTP |
Prometheus Exporter 指标端口 |
四、存储架构
4.1 存储卷清单
| PVC 名称 |
容量 |
存储类 |
绑定节点 |
用途 |
| data-rocketmq-ddffdb1a-0-0-0 |
10Gi |
csi-localpv |
qfusion4 |
Broker-0 数据 |
| data-rocketmq-ddffdb1a-0-1-0 |
10Gi |
csi-localpv |
qfusion1 |
Broker-1 数据 |
| data-rocketmq-ddffdb1a-0-2-0 |
10Gi |
csi-localpv |
qfusion2 |
Broker-2 数据 |
| data-rocketmq-ddffdb1a-nameserver-0-0 |
20Gi |
csi-localpv |
qfusion4 |
NameServer-0 数据 |
| data-rocketmq-ddffdb1a-nameserver-1-0 |
20Gi |
csi-localpv |
qfusion2 |
NameServer-1 数据 |
| data-rocketmq-ddffdb1a-nameserver-2-0 |
20Gi |
csi-localpv |
qfusion1 |
NameServer-2 数据 |
4.2 Broker 存储目录结构
/root/store/
├── commitlog/ # 原始消息存储
├── consumequeue/ # 消费队列索引
├── index/ # 消息索引文件
├── abort/ # 异常退出标识
└── logs/ # 运行日志
├── broker.log # Broker 主日志
├── dledger/ # DLedger 一致性日志
└── rmq.log # 消息追踪日志
4.3 NameServer 存储目录结构
/root/
├── logs/
│ └── namesrv.log # NameServer 日志
└── store/ # NameServer 数据目录(如有)
五、DLedger 模式详解
5.1 DLedger 工作原理
生产者写入流程 (DLedger 模式):
┌────────┐ ┌─────────┐ ┌────────┐ ┌──────────┐ ┌─────────┐
│Producer│ ──> │ Leader │ ──> │ Raft │ ──> │ Followers│ ──> │ 确认写入 │
└────────┘ │ Broker │ │ 协议 │ │ Broker │ └─────────┘
└─────────┘ └──────────┘ └──────────┘
│ │
└────────── 同步复制 ────────────┘
要求: 至少 2/3 节点确认后才返回成功(多数派原则)
5.2 Leader 选举
| 场景 |
行为 |
| Leader 故障 |
剩余 Follower 自动选举新 Leader |
| 网络分区 |
多数派节点组继续服务 |
| 节点恢复 |
自动加入集群成为 Follower |
5.3 DLedger 配置参数
| 参数 |
值 |
说明 |
| clusterMode |
DLEDGER_ASYNC |
DLedger 异步复制模式 |
| replicaPerGroup |
3 |
每组副本数 |
| size |
1 |
Broker 组数 |
六、客户端连接配置
6.1 NameServer 连接地址
# 内部集群访问(Pod 间)
246.108.185.135:9876
# 多地址格式(分号分隔)
245.0.3.219:9876;245.0.1.220:9876;245.0.0.45:9876
6.2 Broker 连接地址
| Broker |
内部地址 |
外部地址 |
| Broker-0 (Leader) |
x.x.x.146:3 |
245.0.3.63:10911 |
| Broker-1 (Follower) |
x.x.x.146:9 |
245.0.0.157:10911 |
| Broker-2 (Follower) |
x.x.x.146:6 |
245.0.1.254:10911 |
6.3 Java 客户端示例
DefaultMQProducer producer = new DefaultMQProducer("producer-group");
producer.setNamesrvAddr("246.108.185.135:9876");
producer.start();
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("consumer-group");
consumer.setNamesrvAddr("246.108.185.135:9876");
consumer.subscribe("bpx-topic", "*");
consumer.registerMessageListener(new MessageListenerConcurrently() {
@Override
public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, Context context) {
return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
}
});
consumer.start();
七、常用资源查询命令
7.1 集群资源查询
export KUBECONFIG=/bpx/.145-admin.conf
export NSRV=246.108.185.135:9876
kubectl get pods -n qfusion-admin | grep rocketmq
kubectl get broker -n qfusion-admin
kubectl get nameservers -n qfusion-admin
kubectl get svc -n qfusion-admin | grep rocketmq
kubectl get pvc -n qfusion-admin | grep rocketmq
7.2 Topic 相关查询
kubectl exec -it -n qfusion-admin rocketmq-ddffdb1a-0-0-0 -- bash
mqadmin topiclist -n $NSRV
mqadmin topicStatus -n $NSRV -t <TOPIC_NAME>
mqadmin topicRoute -n $NSRV -t <TOPIC_NAME>
7.3 消费者相关查询
mqadmin consumerProgress -n $NSRV
mqadmin consumerProgress -n $NSRV -g <CONSUMER_GROUP>
mqadmin consumerConnection -n $NSRV -g <CONSUMER_GROUP>
八、集群模式对比
| 特性 |
DLedger 模式 (当前) |
主从同步模式 |
单机模式 |
| 高可用 |
✅ 自动故障切换 |
✅ 需手动切换 |
❌ |
| 数据一致性 |
✅ 强一致性 (Raft) |
⚠️ 异步复制可能丢失 |
N/A |
| 副本数 |
≥ 3 (推荐) |
1 主 + N 从 |
1 |
| 性能 |
中等 (写入需多数派确认) |
高 |
最高 |
| 运维复杂度 |
中等 |
较高 |
低 |
九、核心概念关系图
┌─────────────────────────────────────────────────────────────────────┐
│ RocketMQ 消息模型 │
├─────────────────────────────────────────────────────────────────────┤
│ │
│ Producer ──发送消息──> Topic ──分片──> Queue ──分布──> Broker │
│ │ │ │
│ │ └──────────────────┐ │
│ │ │ │
│ v v │
│ ConsumerGroup ──订阅──> Topic ──消费──> Queue ──分配──> Consumer │
│ │
│ 关键点: │
│ * 一个 Topic 可以有多个 ConsumerGroup │
│ * 一个 ConsumerGroup 内只有一个 Consumer 消费同一条消息 │
│ * 不同 ConsumerGroup 可以独立消费同一条消息 │
│ * Offset 按 ConsumerGroup 维度管理 │
│ │
└─────────────────────────────────────────────────────────────────────┘
十、快速参考
10.1 集群访问信息
| 项目 |
值 |
| NameServer 地址 |
246.108.185.135:9876 |
| 集群名称 |
rocketmq-ddffdb1a |
| Broker 组名称 |
rocketmq-ddffdb1a-0 |
| 命名空间 |
qfusion-admin |
| Kubeconfig |
/bpx/.145-admin.conf |
| mqadmin 路径 |
/root/rocketmq/broker/bin/mqadmin |
10.2 当前 Topic 示例
| Topic |
类型 |
说明 |
| bpx-topic |
业务 Topic |
测试/业务 Topic |
| %RETRY%bpx-consumer-group |
重试 Topic |
消费失败重试队列 |
| SCHEDULE_TOPIC_XXXX |
延时 Topic |
延时消息专用 |
| RMQ_SYS_TRACE_TOPIC |
系统 Topic |
消息追踪 |
10.3 当前 ConsumerGroup 示例
| ConsumerGroup |
状态 |
说明 |
| bpx-consumer-group |
OFFLINE |
测试消费组(当前离线) |
| TOOLS_CONSUMER |
在线 |
工具消费组(3个消费者) |
文档版本: v1.1
最后更新: 2025-12-31
维护者: 运维团队
十一、入门必读知识
11.1 消息队列是什么?
消息队列(Message Queue)是一种进程间通信或服务间通信的方式:
┌─────────────┐ 消息 ┌─────────────┐ 消息 ┌─────────────┐
│ 服务 A │ ──────────> │ RocketMQ │ ──────────> │ 服务 B │
│ (生产者) │ │ (中转站) │ │ (消费者) │
└─────────────┘ └─────────────┘ └─────────────┘
核心价值:
* 解耦: A 不需要知道 B 的存在
* 异步: A 发完就继续,不用等 B 处理
* 削峰: 高峰期消息先缓存,慢慢消费
11.2 消息的生命周期
1. [生产] Producer 发送消息到 Topic
↓
2. [路由] NameServer 告诉 Producer 该连哪个 Broker
↓
3. [存储] Broker 将消息写入 CommitLog
↓
4. [索引] 更新 ConsumeQueue 索引
↓
5. [推送] Consumer 从 Broker 拉取消息
↓
6. [消费] Consumer 处理业务逻辑
↓
7. [确认] Consumer 提交消费 Offset
↓
8. [清理] 过期消息被自动删除
11.3 消费模式对比
| 模式 |
说明 |
适用场景 |
| 集群模式 (Clustering) |
一个 ConsumerGroup 内每条消息只被一个 Consumer 消费 |
业务处理,避免重复 |
| 广播模式 (Broadcasting) |
每个 Consumer 都会收到所有消息 |
配置下发、通知 |
consumer.setMessageModel(MessageModel.BROADCASTING);
11.4 推送 vs 拉取
| 模式 |
特点 |
代码示例 |
| Push Consumer |
Broker 主动推送,实时性高 |
DefaultMQPushConsumer |
| Pull Consumer |
客户端主动拉取,可控性强 |
DefaultMQPullConsumer |
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("group");
consumer.subscribe("topic", "*");
consumer.registerMessageListener(listener);
consumer.start();
DefaultMQPullConsumer consumer = new DefaultMQPullConsumer("group");
consumer.start();
Set<MessageQueue> mqs = consumer.fetchSubscribeMessageQueues("topic");
11.5 消息重试机制
消费失败后的处理流程:
┌──────────────┐
│ 消息消费失败 │
└──────┬───────┘
│
v
┌──────────────────┐ 是 ┌────────────────┐
│ 是否返回 RECONSUME_LATER? │ ──> │ 加入重试队列 │
└──────────────────┘ └───────┬────────┘
│ 否 │
v v
┌────────────────┐ ┌─────────────────────┐
│ 消息成功确认 │ │ 延时后重新消费 │
│ Offset 推进 │ │ (延时等级递增) │
└────────────────┘ └─────────────────────┘
│
┌──────────┴──────────┐
v v
重试次数 < 16 重试次数 >= 16
│ │
v v
继续重试 进入死信队列
(%RETRY%组名) (%DLQ%组名)
重试延时等级 (默认):
级别 1: 1s 2: 5s 3: 10s 4: 30s 5: 1m
级别 6: 2m 7: 3m 8: 4m 9: 5m 10: 6m
级别 11: 7m 12: 8m 13: 9m 14: 10m 15: 30m
级别 16: 1h
11.6 延时消息
RocketMQ 支持特定延时等级的消息:
Message msg = new Message("topic", "Hello".getBytes());
msg.setDelayTimeLevel(3);
producer.send(msg);
⚠️ 注意: 开源版只支持固定延时等级,不支持任意秒数
11.7 事务消息
事务消息 = 消息发送 + 本地事务 的原子性保证
┌──────────┐ ┌──────────┐ ┌──────────┐
│ 发送半消息 │ ──> │ 执行本地事务 │ ──> │ 提交/回滚 │
└──────────┘ └──────────┘ └──────────┘
│ │
v v
┌──────────────────────────────────────────┐
│ Broker 等待确认 │
│ * 回查机制: 未确认则反查事务状态 │
└──────────────────────────────────────────┘
TransactionMQProducer producer = new TransactionMQProducer("group");
producer.setTransactionListener(new TransactionListener() {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
return LocalTransactionState.COMMIT_MESSAGE;
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
return LocalTransactionState.COMMIT_MESSAGE;
}
});
11.8 Tag 与 Message Key
Topic: 第一级分类,如 "订单"
Tag: 第二级分类,如 "创建"、"支付"、"完成"
Key: 消息唯一标识,如订单ID
┌─────────────────────────────────────────────────┐
│ Topic: OrderTopic │
├─────────────────────────────────────────────────┤
│ Tag: TagA Tag: TagB Tag: TagC │
│ [订单创建] [订单支付] [订单完成] │
│ Key: Order123 Key: Order123 Key: Order123│
└─────────────────────────────────────────────────┘
Message msg = new Message("OrderTopic", "TagA", "Order123", body);
producer.send(msg);
consumer.subscribe("OrderTopic", "TagA || TagB");
11.9 消息堆积与处理
消息堆积 = 生产速度 > 消费速度
┌────────────────────────────────────────────────────────┐
│ Broker │
│ ┌────┐┌────┐┌────┐┌────┐┌────┐┌────┐┌────┐┌────┐ │
│ │msg1││msg2││msg3││msg4││msg5││msg6││msg7││msg8│...│
│ └────┘└────┘└────┘└────┘└────┘└────┘└────┘└────┘ │
│ ↑ ↑ │
│ 已消费 堆积部分 │
│ │ │
│ Diff Total = 堆积量 │
└──────────────────────────────────────────────────────┘
堆积处理策略:
| 策略 |
方法 |
风险 |
| 扩容消费者 |
增加 Consumer 数量 |
无风险 |
| 跳过堆积 |
skipAccumulatedMessage |
丢失消息 |
| 重置 Offset |
resetOffsetByTime |
重复消费或丢失 |
11.10 常见问题 FAQ
| 问题 |
可能原因 |
排查方法 |
| 消息发送失败 |
NameServer 连不上 |
检查网络、防火墙 |
| 消费不到消息 |
订阅关系错误 |
检查 Tag 表达式 |
| 消息重复消费 |
Offset 提交失败 |
实现幂等性 |
| 消息丢失 |
Broker 故障 |
检查集群健康度 |
| 堆积持续增长 |
消费者处理慢 |
增加消费者或优化处理逻辑 |
| 无法连接 Broker |
防火墙/端口 |
开放 10911/9876 端口 |
十二、学习路径建议
12.1 新手入门路线
第1天: 基础概念
├─ 了解消息队列是什么
├─ 理解 Topic/ConsumerGroup/Queue 概念
└─ 熟悉 NameServer 和 Broker 角色
第2天: 环境操作
├─ 使用 mqadmin 查看集群状态
├─ 查看 Topic 和 ConsumerGroup
└─ 理解 Offset 和消息堆积
第3天: 生产消费
├─ 使用 mqadmin sendMessage 发送消息
├─ 使用 mqadmin consumeMessage 消费消息
└─ 观察 topicStatus 变化
第4天: 故障排查
├─ 掌握 4 步排障法
├─ 理解重试机制
└─ 学习处理堆积
第5天: 进阶知识
├─ DLedger 一致性原理
├─ 事务消息机制
└─ 性能调优基础
12.2 实战练习建议
kubectl exec -n qfusion-admin rocketmq-ddffdb1a-0-0-0 -- \
/root/rocketmq/broker/bin/mqadmin sendMessage \
-n 246.108.185.135:9876 -t bpx-topic -p "test"
kubectl exec -n qfusion-admin rocketmq-ddffdb1a-0-0-0 -- \
/root/rocketmq/broker/bin/mqadmin topicStatus \
-n 246.108.185.135:9876 -t bpx-topic
kubectl exec -n qfusion-admin rocketmq-ddffdb1a-0-0-0 -- \
/root/rocketmq/broker/bin/mqadmin consumerProgress \
-n 246.108.185.135:9876
kubectl exec -n qfusion-admin rocketmq-ddffdb1a-0-0-0 -- \
/root/rocketmq/broker/bin/mqadmin statsAll \
-n 246.108.185.135:9876 | grep bpx-topic
12.3 推荐阅读顺序
- 本文档 - 架构与术语
- mqadmin-troubleshooting-guide.md - 命令工具使用
- RocketMQ_Emergency_Troubleshooting.md - 应急故障处理
- 官方文档 - https://rocketmq.apache.org/zh/docs/
所有评论(0)