RocketMQ 深度解析
#组件 · #RocketMQ · #消息队列 · #事务消息 · #阿里 · #分布式
阿里开源的万亿级消息引擎——高可靠存储、分布式事务、延迟消息与金融级特性。
一、架构
1.1 四大角色
mermaid
flowchart TB
subgraph NS["NameServer Cluster<br/>无状态, 路由注册与发现"]
NS1["NameServer 1"]
NS2["NameServer 2"]
end
subgraph Brokers["Broker 集群"]
subgraph B1["Broker 1"]
BM1["Master (写)"]
BS1["Slave (备, 读)"]
end
subgraph B2["Broker 2"]
BM2["Master (写)"]
BS2["Slave (备, 读)"]
end
end
Brokers -->|"定时注册<br/>Topic → Broker 路由"| NS
Producer["Producer<br/>生产者"] -->|"从 NameServer<br/>获取路由"| NS
Producer -->|"发送消息"| Brokers
Consumer["Consumer<br/>消费者"] -->|"从 NameServer<br/>获取路由"| NS
Brokers -->|"拉取消息 (Pull)"| Consumer与 Kafka 关键区别:Broker 主从靠手动/自动切换,没有 ZK/KRaft 选举 Leader。高可用靠 DLedger(Raft 变体,4.5+)或 Controller 模式(5.0+)。
| 角色 | 职责 |
|---|---|
| NameServer | 路由注册中心 (Topic → Broker 映射),无状态,互相不通信 |
| Broker | 消息存储与投递,主从模式,定时向所有 NameServer 注册 |
| Producer | 从 NameServer 拉路由 → 发送消息到 Broker |
| Consumer | 从 NameServer 拉路由 → 拉取消息(Pull 模型) |
与 Kafka 架构关键不同:Broker 主从依赖手动/自动切换,没有 ZK/KRaft 选举 Leader。高可用靠 DLedger(Raft 变体,4.5+)或 Controller 模式(5.0+)。
1.2 Broker 集群模式
| 模式 | 特点 |
|---|---|
| 单 Master | 测试用,宕机不可用 |
| 多 Master | 无 Slave,宕机会丢少量消息 |
| 多 Master 多 Slave (异步) | 主挂备升,可能丢少量 |
| 多 Master 多 Slave (同步双写) | 主备同时写,不丢消息但性能下降 |
| DLedger (Raft) | 自动选主,强一致,金融级推荐 |
二、消息存储(CommitLog + ConsumeQueue)
2.1 存储架构
RocketMQ 文件布局 (单 Broker):
store/
├── commitlog/ ← 所有 Topic 消息顺序追加写入
│ ├── 00000000000000000000 (1GB/文件)
│ ├── 00000000001073741824
│ └── ...
├── consumequeue/ ← 按 Topic-Queue 的索引
│ └── TopicA/
│ ├── 0/00000000000000000000
│ └── 1/00000000000000000000
├── index/ ← 按 Key/时间查询的哈希索引
│ └── ...
└── config/2.2 写入与消费流程
Producer 发送消息:
消息 → CommitLog (顺序写, 单文件, 1GB轮转)
↓ (后台 ReputService 异步构建)
ConsumeQueue (Topic-Queue 粒度索引)
↓ (消费者拉取)
Consumer
ConsumeQueue 条目 (定长 20B):
┌────────────┬──────────┬──────────┐
│ CommitLog │ 消息大小 │ Tag Hash │
│ Offset (8B)│ (4B) │ (8B) │
└────────────┴──────────┴──────────┘优势:所有 Topic 共享一个 CommitLog,顺序写只需一个文件,写入极快。
对比 Kafka:Kafka 按 Partition 单独写文件,Partition 多时随机写增多;RocketMQ 所有消息合并顺序写,仅需少量随机读。
2.3 刷盘与 HA
| 参数 | 说明 |
|---|---|
flushDiskType=ASYNC_FLUSH | 异步刷盘(默认,吞吐高) |
flushDiskType=SYNC_FLUSH | 同步刷盘(每条消息等 fsync) |
brokerRole=SYNC_MASTER | 同步双写 Slave(主等备确认) |
可靠级别:异步刷盘+异步复制 < 同步复制 < 同步刷盘+同步复制(最安全但最低)
三、消息模型
3.1 消费模式
集群消费 (默认, CLUSTERING):
Topic: Order
Queue 0 → Consumer A1 ─┐
Queue 1 → Consumer A2 ├ Consumer Group A
Queue 2 → Consumer A3 ─┘
(每条消息只被组内一个消费者消费)
广播消费 (BROADCASTING):
Topic: Notify
Queue 0 → Consumer B1, Consumer B2, Consumer B3
Queue 1 → Consumer B1, Consumer B2, Consumer B3
(每条消息被组内所有消费者都收到)3.2 消费进度(Offset)
Consumer Offset 存储在 Broker (非 ZK):
集群模式: 存 Broker 端
广播模式: 存 Consumer 本地
消费进度: minOffset ←── 已消费 ──→ consumerOffset ←── 未消费 ──→ maxOffset
重新消费: 重置 consumerOffset 到某个时间点3.3 Rebalance
触发条件:
- 消费者数量变化 (上下线)
- Topic Queue 数量变化
- 消费者心跳超时
策略 (5.0 统一为兼容 Kafka 的协作式):
AllocateMessageQueueAveragely ← 平均分配 (默认)
AllocateMessageQueueAveragelyByCircle ← 环形平均
AllocateMessageQueueConsistentHash ← 一致性 Hash
AllocateMachineRoomNearby ← 同机房优先四、事务消息
4.1 分布式事务流程
Producer Broker Consumer
│ │ │
│ 1. send half msg │ │
│─────────────────────→│ │
│ (半消息, 对消费者不可见) │
│ │ │
│ 2. 执行本地事务 │ │
│ (扣款/减库存...) │ │
│ │ │
│ 3. commit / rollback │ │
│─────────────────────→│ │
│ │ │
│ ┌──────┴──────┐ │
│ │ commit? │ │
│ │ 将半消息标记 │ │
│ │ 为可投递 │ │
│ └──────┬──────┘ │
│ │ │
│ │ 4. 投递 │
│ │─────────────────────→│
│ │ │
│ 如果步骤 3 超时 │ │
│ Broker 回查事务状态 │ │
│ ←───────────────────→│ │4.2 使用示例
java
// 发送事务消息
TransactionMQProducer producer = new TransactionMQProducer("tx_group");
producer.setTransactionListener(new TransactionListener() {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// 执行本地事务
try {
orderService.createOrder(); // 本地 DB 操作
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// Broker 回查: 检查本地事务是否提交
String orderId = msg.getKeys();
return orderService.isOrderExists(orderId)
? LocalTransactionState.COMMIT_MESSAGE
: LocalTransactionState.ROLLBACK_MESSAGE;
}
});
producer.start();
// 发送
Message msg = new Message("OrderTopic", "order.create", orderId, body);
producer.sendMessageInTransaction(msg, null);五、延迟消息
RocketMQ 开箱即用,支持 18 个等级:
1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2hjava
Message msg = new Message("Topic", "Tag", body);
msg.setDelayTimeLevel(3); // 等级 3 = 10s
producer.send(msg);原理:延迟消息先存在 SCHEDULE_TOPIC_XXXX 内部 Topic,定时任务到点捞取投递到目标 Topic。
六、消息过滤
| 方式 | 说明 |
|---|---|
| Tag 过滤 | 消息标记 Tag,消费端订阅 TopicA || Tag1 || Tag2 |
| SQL92 过滤 | 按消息属性 a > 5 AND b = 'abc'(需 Broker 开启) |
| 类过滤 (4.x) | Java 序列化过滤(不推荐) |
七、性能与最佳实践
| 优化 | 说明 |
|---|---|
| 消息大小 | < 512KB(默认 4MB 限制) |
| Queue 数量 | 与消费者数量匹配,不宜过多(影响 CommitLog 随机读) |
| 异步发送 | 吞吐提升 3-5 倍 |
| 批量消费 | consumeMessageBatchMaxSize 一次拉多条 |
| 顺序消息 | 同一业务 ID 发到同一 Queue(MessageQueueSelector) |
顺序消息保证
java
// 全局有序: 一个 Topic 只有一个 Queue
// 分区有序 (常见):
producer.send(msg, new MessageQueueSelector() {
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
long orderId = (Long) arg;
int index = (int) (orderId % mqs.size());
return mqs.get(index);
}
}, orderId);八、Kafka vs RocketMQ 技术选型
| 维度 | Kafka | RocketMQ |
|---|---|---|
| 事务消息 | ✅ Kafka Streams 事务 | ✅ 原生事务消息(半消息+回查) |
| 延迟消息 | ❌ 需手动实现 | ✅ 原生 18 级延迟 |
| 顺序消息 | 分区内有序 | Queue 内有序 |
| 消息重试 | 需自行死信处理 | ✅ 原生重试 + 死信队列 |
| 堆积能力 | 强(千亿级) | 强(百亿级) |
| 吞吐量 | 极高 | 高 |
| 运维复杂度 | 中等 | 较高 |
| 生态 | 大数据/流处理 | 阿里云商业版、Spring Cloud Stream |
| 适用 | 日志、流、大数据 | 电商、金融、分布式事务 |
参考
九、RocketMQ 5.0 新特性
架构演进
mermaid
flowchart TB
subgraph V4["RocketMQ 4.x"]
NS4["NameServer<br/>(路由)"]
B4["Broker<br/>(存储+计算)"]
P4["Producer"]
C4["Consumer"]
P4 --> B4
B4 --> C4
B4 --> NS4
end
subgraph V5["RocketMQ 5.0"]
NS5["NameServer<br/>(路由)"]
Proxy["Proxy<br/>(接入层, gRPC)"]
B5["Broker<br/>(存储)"]
Controller["Controller<br/>(选主, Raft)"]
P5["Producer<br/>(gRPC SDK)"]
C5["Consumer<br/>(Pop 模式)"]
P5 -->|"gRPC"| Proxy
Proxy --> B5
B5 --> Controller
Proxy --> C5
end5.0 核心变化
| 特性 | 4.x | 5.0 | 说明 |
|---|---|---|---|
| 接入协议 | 自定义 Remoting | gRPC | 多语言 SDK 统一,跨语言更方便 |
| 消费模式 | Pull + Push | Pop 消费 | 无状态消费,无需 Rebalance |
| 高可用 | DLedger (可选) | Controller 模式 | Raft 选主,自动故障转移 |
| 延迟消息 | 18 个固定级别 | 任意时间延迟 | 基于时间轮,精度秒级 |
| 架构 | Broker 耦合 | Proxy + Broker 分离 | 计算存储分离,独立扩展 |
| 可观测性 | 基础 | OpenTelemetry | 原生链路追踪 |
Pop 消费模式
传统 Pull 模式问题:
- Consumer 与 Queue 绑定 → Rebalance 时消费暂停
- Consumer 挂了 → 该 Queue 消息积压直到 Rebalance 完成
Pop 模式(5.0 新增):
- Consumer 无状态,每次 Pop 请求从 Broker 获取消息
- 无需 Rebalance,Consumer 随时加入/退出
- Broker 端管理消费进度(类似 RabbitMQ 的 Push)
- 适合 Serverless / 弹性伸缩场景java
// 5.0 gRPC SDK 示例
ClientServiceProvider provider = ClientServiceProvider.loadService();
ClientConfiguration configuration = ClientConfiguration.newBuilder()
.setEndpoints("proxy-host:8081")
.build();
// Simple Consumer(Pop 模式)
SimpleConsumer consumer = provider.newSimpleConsumerBuilder()
.setClientConfiguration(configuration)
.setConsumerGroup("my-group")
.setSubscriptionExpressions(Collections.singletonMap("TopicA", FilterExpression.SUB_ALL))
.build();
// 主动拉取(无 Rebalance)
List<MessageView> messages = consumer.receive(32, Duration.ofSeconds(10));
for (MessageView msg : messages) {
// 处理消息
consumer.ack(msg); // 确认
}十、运维排查
常用运维命令
bash
# 集群状态
mqadmin clusterList -n localhost:9876
# Topic 列表
mqadmin topicList -n localhost:9876
# Topic 详情(Queue 分布)
mqadmin topicStatus -n localhost:9876 -t OrderTopic
# 消费者组状态(消费进度、积压)
mqadmin consumerProgress -n localhost:9876 -g my-consumer-group
# 查看消息积压(最重要的运维指标)
mqadmin consumerProgress -n localhost:9876 -g my-group
# 输出:
# Topic Broker QueueId Broker Offset Consumer Offset Diff
# OrderTopic broker-a 0 100000 99500 500
# OrderTopic broker-a 1 100000 98000 2000 ← 积压!
# 按 MessageId 查询消息
mqadmin queryMsgById -n localhost:9876 -i 0A0A0A0A00002A9F000000000000001E
# 按 Key 查询消息
mqadmin queryMsgByKey -n localhost:9876 -t OrderTopic -k order_12345
# 按时间范围查询
mqadmin queryMsgByOffset -n localhost:9876 -t OrderTopic -b broker-a -i 0 -o 99000
# 重置消费位点(回溯消费)
mqadmin resetOffsetByTime -n localhost:9876 -g my-group -t OrderTopic -s "2024-01-01#00:00:00:000"
# 查看 Broker 运行状态
mqadmin brokerStatus -n localhost:9876 -b broker-a:10911常见问题排查
| 问题 | 排查命令 | 解决方案 |
|---|---|---|
| 消息积压 | consumerProgress 看 Diff | 增加消费者 / 增加线程 / 跳过非关键消息 |
| 发送失败 | 查看 Producer 日志 + brokerStatus | 检查 Broker 磁盘/内存/网络 |
| 消费重复 | 查看 Consumer 日志 + Rebalance 频率 | 实现幂等 / 增大 maxReconsumeTimes |
| Broker 主从切换 | clusterList 看 Master 变化 | 检查 DLedger/Controller 日志 |
| 消息丢失 | 对比 Producer 发送数和 Consumer 消费数 | 开启同步刷盘 + 同步复制 |
| 延迟消息不准 | 检查 SCHEDULE_TOPIC_XXXX 积压 | Broker 负载过高,扩容 |
监控指标
promql
# RocketMQ Exporter 关键指标
# 消费积压(最重要)
rocketmq_consumer_diff{group="my-group", topic="OrderTopic"}
# Broker 写入 TPS
rocketmq_broker_tps{broker="broker-a"}
# 发送延迟
rocketmq_producer_send_latency_bucket
# 消费延迟
rocketmq_consumer_consume_latency_bucket
# 告警规则
- alert: RocketMQConsumerLag
expr: rocketmq_consumer_diff > 10000
for: 5m
labels:
severity: warning
annotations:
summary: "RocketMQ 消费积压 {{ $value }} 条"
登录后即可发表评论 👇