主题
RocketMQ 详解:架构、事务消息与延时消息
本文是消息队列系统学习系列的 L2 核心篇。前置:28. kafka-producer-consumer-internals、29. kafka-reliability-high-performance。 学完可以配合面试题食用:13-rocketmq-architecture-nameserver-broker-producer-consumer、14-rocketmq-transaction-message、16-delayed-message-solutions-comparison
为什么需要再学一个 MQ
Kafka 已经是日志流的事实标准了,但落到业务消息场景——事务一致性、任意延时、精确的消息轨迹——Kafka 的原生能力就很勉强。事务只保跨分区的 exactly-once 写入,不保跨系统事务;延时消息要自己造轮子。
RocketMQ 是阿里在 Kafka 思路上的批量改造产物,保留了 Kafka 的日志存储和顺序写,但把消息模型从"拉日志"拉回"发消息"。它对业务场景的黑盒需求(事务、延时、消息轨迹、死信重试)做了内建支持,这也是为什么国内电商、支付、直播场景大量用 RocketMQ 而不是 Kafka。
架构:NameServer + Broker + 更轻的协调
RocketMQ 的架构比 Kafka 少一个依赖——没有 ZooKeeper,用 NameServer 替代。
mermaid
graph LR
subgraph Producer
P1[Producer Cluster]
end
subgraph NameServer
NS1[NameServer 1]
NS2[NameServer 2]
end
subgraph Broker
B1[Broker A - Master]
B2[Broker A - Slave]
B3[Broker B - Master]
B4[Broker B - Slave]
end
subgraph Consumer
C1[Consumer Group]
end
P1 -->|注册/路由查询| NS1
P1 -->|注册/路由查询| NS2
B1 -->|心跳注册| NS1
B1 -->|心跳注册| NS2
B3 -->|心跳注册| NS1
B3 -->|心跳注册| NS2
NS1 -->|Topic 路由| C1
NS2 -->|Topic 路由| C1
P1 -->|发消息| B1
P1 -->|发消息| B3
B1 -->|主从复制| B2
B3 -->|主从复制| B4
C1 -->|拉消息| B1
C1 -->|拉消息| B3NameServer 做了什么? 每个 NameServer 独立维护全量路由表,无状态。Broker 启动后向所有 NameServer 注册 Topic 路由(Broker 地址 + 队列数),每 30 秒心跳续命。Producer/Consumer 启动时从 NameServer 拉路由,缓存在本地。NameServer 挂一个不影响集群,因为客户端走的是本地缓存,且会尝试其他 NameServer。
和 Kafka 有什么不同?
| 维度 | Kafka | RocketMQ |
|---|---|---|
| 元数据服务 | ZooKeeper / KRaft | NameServer(无状态) |
| 存储单元 | Partition → Segment | ConsumeQueue + CommitLog |
| 队列 vs 分区 | 无 Queue 概念 | 每个 Topic 下有多个 Queue |
| 消费偏移 | 自动提交到 __consumer_offsets | 消费进度由 Consumer 自己管理 |
| 消费模式 | 分区内有序 | 支持集群/广播两种模式 |
RocketMQ 的 CommitLog 是唯一的顺序写文件,ConsumeQueue 只是 CommitLog 的索引——这和 Kafka 每个 Partition 独立 Segment 的设计不同,让 RocketMQ 的单机文件数可控,但队列之间的隔离性不如 Kafka。
事务消息:半消息 + 本地事务 + 回查
事务消息是 RocketMQ 最出圈的特性。它的核心矛盾是"发送消息"和"本地 DB 事务"必须捆绑一致——要么都成功,要么都回滚。
RocketMQ 的解法是两阶段 + 事务回查:
mermaid
sequenceDiagram
participant App as 应用
participant MQ as Broker
App->>MQ: 1. 发送半消息(prepare)
MQ-->>App: 半消息 OK
App->>App: 2. 执行本地事务(DB 操作)
alt 本地事务成功
App->>MQ: 3. commit
MQ-->>Consumer: 投递消息
else 本地事务失败
App->>MQ: 3. rollback
MQ->>MQ: 删除半消息
else 超时未响应
MQ->>App: 4. 回查事务状态
App-->>MQ: 返回 commit/rollback
end看看代码怎么落地:
java
// 1. 生产者端:设置事务监听器
TransactionMQProducer producer = new TransactionMQProducer("tx-group");
producer.setNamesrvAddr("localhost:9876");
producer.setTransactionListener(new TransactionListener() {
@Override
public LocalTransactionState executeLocalTransaction(Message msg, Object arg) {
// 参数 arg 是业务参数,比如订单对象
Order order = (Order) arg;
try {
// 执行业务 DB 操作
orderService.create(order); // insert into orders
return LocalTransactionState.COMMIT_MESSAGE;
} catch (Exception e) {
return LocalTransactionState.ROLLBACK_MESSAGE;
}
}
@Override
public LocalTransactionState checkLocalTransaction(MessageExt msg) {
// Broker 回查:根据消息体中的业务 key 查 DB
String orderId = msg.getKeys();
// select * from orders where id = orderId
if (orderService.exists(orderId)) {
return LocalTransactionState.COMMIT_MESSAGE;
}
return LocalTransactionState.ROLLBACK_MESSAGE;
}
});
producer.start();
// 2. 发送半消息
Message msg = new Message("order-tx", "create".getBytes());
msg.setKeys(orderId); // 回查时用
SendResult result = producer.sendMessageInTransaction(msg, order);几个关键点:
- 半消息对 Consumer 不可见,只有 commit 后 Consumer 才收到
- 回查默认间隔 60 秒,可配置通过
transactionCheckInterval - 回查次数默认 15 次,超限自动 rollback
- 事务回查要保证幂等——同一个 orderId 查 DB 多次返回相同结果
延时消息:18 个固定级别,够用但不灵活
RocketMQ 的延时消息不是任意延时,而是 18 个固定级别:
java
// 1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
// 对应 level = 1 ~ 18
Message msg = new Message("delayed-topic", "payload".getBytes());
msg.setDelayTimeLevel(5); // 延时 1 分钟
producer.send(msg);实现原理: Broker 启动时对每个延时级别创建一个 SCHEDULE_TOPIC_XXXX 队列。消息进来时,先写入对应的 SCHEDULE_TOPIC 队列,Broker 内部有一个定时任务(SchedulerService),每秒扫描到期消息,转移到目标 Topic 的队列。这个机制叫定时转存,精度在 1 秒左右,不会有 Kafka 那种写完到期再消费的延迟累积。
为什么不支持任意延时? 每个级别一个队列,固定数量好管理;如果支持任意延时,需要维护一个巨大的定时器数据结构(如时间轮),RocketMQ 4.x 选了简单实现。5.x 版本开始支持精确延时,但生产环境仍然以固定级别为主。
如果业务要 13 秒延时怎么办? 两种方案:
- 取 10s(level 4),应用层做 3 秒 Timer 补偿
- 自己用时间轮实现:在内存里维护一个延时队列,到期后发 RocketMQ 消息
为什么国内电商多用 RocketMQ
RocketMQ 在阿里内部的业务场景打磨出来的特性,刚好是电商的刚需:
- 事务消息:下单扣库存 + 发消息必须原子
- 延时消息:订单超时关闭、支付结果轮询间隔
- 消息轨迹:一条消息从生产到消费的完整链路,debug 不需要翻日志
- 消息重试与死信:消费失败自动重试 16 次,超限进死信队列,控制台可以直接重置
- 消费端 pull 长轮询:消费者拉不到消息时,Broker 会 hold 住请求最多 15 秒,消息到了立即返回,减少空轮询
Kafka 在吞吐量(单机百万 TPS)上仍然碾压 RocketMQ(单机十万级),但 RocketMQ 在业务消息的完整度上更胜一筹。选型不是哪个更好,而是你要的是"日志管道"还是"业务消息总线"。
常见误区与小结
- 误区:事务消息就是 XA 分布式事务。 不是。事务消息只保证"发送消息"和"本地事务"一致,不保证消费端成功。如果消费端失败,需要重试/死信机制兜底,这是两回事。
- 误区:延时消息精度很高。 固定级别延时只有秒级精度,大量 13 秒、17 秒的延时需求需要业务层补偿或自己实现。
- 误区:NameServer 挂了整个集群就完了。 客户端有本地缓存,NameServer 瞬时不可用不影响已有连接,但新 Topic 路由无法获取,长期不可用需重启客户端。
- 误区:RocketMQ 比 Kafka 吞吐差很多就不能用。 单机十万级 TPS 对绝大多数业务场景已经够用,只有日志采集、埋点、IoT 数据流这种场景才需要 Kafka 的百万级吞吐。
小结: RocketMQ 用 NameServer 取代了 ZooKeeper,用事务消息和延时消息覆盖了 Kafka 的业务空白。它的核心哲学是"把业务消息的黑盒需求做进 Broker 而非堆在应用层"。下一篇 31. mq-selection-capacity-planning 会从选型决策树出发,把 Kafka、RocketMQ、RabbitMQ、Pulsar 放在一起对比,给出具体场景的选型结论。
参考
- RocketMQ 官方文档:事务消息
- RocketMQ 源码:TransactionMessageBridge — 半消息处理核心类
- 《RocketMQ 技术内幕》丁威 — 第 4 章事务消息、第 6 章延时消息实现细节