消息队列02:RocketMQ
从 NameServer 路由、Broker 存储模型、事务消息两阶段提交到顺序消息实现,拆解 RocketMQ 在业务消息场景下的设计逻辑和可靠性保证。
RocketMQ 不是为”日志管道”设计的,而是为金融级业务消息设计的。它的核心场景是订单、支付、库存同步这类”不能丢消息、要保证顺序、需要事务联动”的链路。
这篇文章围绕四个核心机制:NameServer、Broker 存储模型、事务消息、顺序消息,把 RocketMQ 与 Kafka 的差异和业务适配讲清楚。
为什么业务消息选 RocketMQ
先定位清楚:RocketMQ 的核心场景是金融级可靠事务消息,追求消息零丢失和严格顺序,适合订单、交易、库存等核心业务链路。
与 Kafka 的核心差异:
| 维度 | RocketMQ | Kafka |
|---|---|---|
| 存储模型 | 集中式 CommitLog + ConsumeQueue 索引 | 分区式存储,每 Partition 独立文件 |
| Topic 数量敏感度 | 单机支持上千 Topic 无明显衰减 | Topic > 50 时 PageCache 污染严重 |
| 消息可靠性 | 同步刷盘 + 同步复制,业务消息可靠性更强 | 依赖 ISR + acks=all,需调优 |
| 事务消息 | 原生事务消息 + 事务回查,贴近交易场景 | 跨 Partition 原子写入,缺乏服务端回查 |
| 延迟消息 | 原生支持 18 个延迟级别 | 不原生支持,需额外实现 |
| 死信队列 | 原生死信队列,消费失败自动进入 | 需自己实现 |
简单判断:订单、支付、交易链路、需要延迟消息、Topic 数量多 → 选 RocketMQ。日志采集、实时计算、埋点流处理 → 选 Kafka。
核心架构:NameServer、Broker、Producer、Consumer
RocketMQ 的架构比 Kafka 更简单,但职责清晰:
- NameServer:路由注册中心,类似注册中心,Broker 向 NameServer 注册路由信息,Producer/Consumer 从 NameServer 获取 Broker 地址
- Broker:消息服务器,负责存储、投递和查询,支持主从架构和 DLedger 集群
- Producer:消息生产者,支持同步发送、异步发送、单向发送和事务消息
- Consumer:消息消费者,支持 Push 模式(长轮询)和 Pull 模式
核心设计思路:集中式存储做高可靠性,NameServer 做轻量路由,长轮询做实时消费。
NameServer 与 Kafka 的 ZooKeeper/KRaft 不同:NameServer 不做选主,只做路由注册,每个 NameServer 独立运行,Producer/Consumer 可以连接任意 NameServer。Broker 定期向所有 NameServer 注册路由,NameServer 无状态,故障不影响消息投递(Producer 可以切换到其他 NameServer)。
存储模型:CommitLog + ConsumeQueue
RocketMQ 的存储模型是与 Kafka 最大的差异点。
CommitLog:全局顺序写
所有 Topic 的消息混写在一个全局的 CommitLog 文件中,顺序追加写入:
- 顺序写 I/O 极高:不管有多少 Topic,写入都是顺序写
- Topic 数量不敏感:不像 Kafka,Topic 过多时 Partition 文件分散导致随机写
- 文件结构:CommitLog 按固定大小切分(默认 1GB),文件名是起始 Offset
CommitLog 是物理存储,消息写入后不会删除(只过期清理),消息内容按顺序排列,每条消息带 Topic、QueueId、消息体、属性等。
ConsumeQueue:逻辑索引
每个 Topic 下有多个 Queue(类似 Kafka 的 Partition),每个 Queue 对应一个 ConsumeQueue:
- ConsumeQueue 是索引文件:存储 CommitLog 中的消息偏移量、消息大小和 Tags Hash
- 消费时二次寻址:先读 ConsumeQueue 获取偏移量,再读 CommitLog 获取消息内容
- 写入成本低:Queue 数量增多不会导致 I/O 抖动(只是索引文件增多)
这种设计的取舍:
| 维度 | RocketMQ(CommitLog + ConsumeQueue) | Kafka(Partition 分段) |
|---|---|---|
| 写入性能 | 极高,所有 Topic 混写,顺序 I/O | Partition 过多时退化为随机写 |
| Topic 数量 | 单机支持上千 Topic 无明显衰减 | Topic 过多时 PageCache 污染 |
| 消费性能 | 二次寻址,略低 | 直接读 Partition,顺序读局部性好 |
| 适用场景 | 业务消息,Topic 多,可靠性优先 | 日志流,吞吐优先 |
刷盘策略:同步刷盘 vs 异步刷盘
Broker 的刷盘策略决定消息落盘的可靠性:
| 刺盘策略 | 说明 | 可靠性 | 性能 |
|---|---|---|---|
| 同步刷盘(SYNC_FLUSH) | 消息写入 PageCache 后,立即刷到磁盘才返回成功 | 最高,磁盘故障前不丢消息 | 较低 |
| 异步刷盘(ASYNC_FLUSH) | 消息写入 PageCache 后立即返回,后台定期刷盘 | 较低,Broker 故障可能丢 PageCache 中的消息 | 较高 |
业务消息推荐同步刷盘,日志场景可以用异步刷盘。
主从复制:同步复制 vs 异步复制
Broker 的主从复制策略决定主备一致性:
| 复制策略 | 说明 | 可靠性 | 性能 |
|---|---|---|---|
| 同步复制(SYNC_MASTER) | Master 收到消息后,等待 Slave 同步成功才返回 | 最高,主备一致 | 较低 |
| 异步复制(ASYNC_MASTER) | Master 写入成功后立即返回,Slave 异步同步 | 较低,Master 故障可能丢未同步的消息 | 较高 |
金融场景推荐同步刷盘 + 同步复制,日志场景可以用异步刷盘 + 异步复制。
事务消息:两阶段提交 + 事务回查
RocketMQ 的事务消息是它与 Kafka 最大的差异化能力:保证本地事务和消息发送要么同时成功,要么同时失败。
事务消息流程
事务消息采用两阶段提交 + 事务状态回查:
- 发送 Half 消息:Producer 发送半消息(Half Message)到 Broker,消息对消费者不可见,Broker 返回发送成功
- 执行本地事务:Producer 执行本地业务逻辑(如写入数据库)
- 提交或回滚:
- 本地事务成功 → Producer 发送 Commit 消息,Broker 将 Half 消息变为可消费
- 本地事务失败 → Producer 发送 Rollback 消息,Broker 删除 Half 消息
- 事务回查:如果 Producer 未返回 Commit/Rollback(如网络中断),Broker 定期回查 Producer 的本地事务状态,Producer 根据本地事务表返回结果
为什么需要事务回查
Producer 发送 Half 消息后,网络中断导致 Commit/Rollback 未送达,Broker 不知道该消息应该投递还是删除。事务回查机制让 Broker 主动询问 Producer 的本地事务状态,直到事务结束。
这要求 Producer:
- 维护本地事务状态表(记录 Half 消息 ID 和事务结果)
- 提供事务回查接口(Broker 回调时查询状态表)
与 Kafka 事务的差异
| 维度 | RocketMQ 事务消息 | Kafka 事务 |
|---|---|---|
| 目标 | 本地事务 + 消息发送的一致性 | 跨 Partition 原子写入 |
| 适用场景 | 订单创建 + 消息通知、支付成功 + 状态同步 | 消费-处理-生产链路(流处理) |
| 服务端回查 | 支持,Producer 故障恢复后 Broker 主动回查 | 不支持,需要 Producer 自己处理 |
| 与外部 DB 联动 | 可以,本地事务是 DB 操作,消息发送是 RocketMQ | 不适合,Kafka 事务不覆盖外部 DB |
典型场景:用户下单后,本地写入订单表,同时发送消息通知库存、积分、物流。如果本地事务失败,消息不投递;如果消息发送失败,本地事务回滚。
顺序消息:Queue 级有序
RocketMQ 的顺序消息是 Queue 级有序,不是全局有序。
顺序消息的实现
发送顺序消息时,Producer 使用 MessageQueueSelector 选择同一个 Queue:
producer.send(message, new MessageQueueSelector() {
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
Long orderId = (Long) arg; // 订单 ID
int index = (int) (orderId % mqs.size());
return mqs.get(index);
}
}, orderId);
同一个订单 ID 的消息始终发到同一个 Queue,Queue 内消息有序。
消费顺序消息时,Consumer 使用 MessageListenerOrderly:
consumer.registerMessageListener(new MessageListenerOrderly() {
@Override
public ConsumeOrderlyStatus consumeMessage(List<MessageExt> msgs, ConsumeOrderlyContext context) {
// 单线程按顺序消费
for (MessageExt msg : msgs) {
// 处理消息
}
return ConsumeOrderlyStatus.SUCCESS;
}
});
MessageListenerOrderly 会加锁(synchronized) 并维护全局消费进度,确保同一 Queue 的消息逐条被消费。
顺序消息的代价
| 维度 | 顺序消息 | 并发消息 |
|---|---|---|
| 吞吐量 | 较低,单 Queue 单线程消费 | 较高,多 Queue 多线程并发 |
| 并行度 | Queue 数量 = 最大并行度 | Consumer 线程池并行 |
| 适用场景 | 订单状态变更、库存扣减、支付状态同步 | 日志采集、埋点上报、通知广播 |
顺序消息的并行度受限于 Queue 数量。如果业务要求严格有序,只能单 Queue 单线程,吞吐会很低。大多数场景只需要”局部有序”(同一订单的消息有序),可以通过订单 ID 分 Queue 实现。
延迟消息:18 个延迟级别
RocketMQ 原生支持延迟消息,不需要额外实现时间轮。内置 18 个延迟级别:
| 级别 | 延迟时间 |
|---|---|
| 1 | 1s |
| 2 | 5s |
| 3 | 10s |
| 4 | 30s |
| 5 | 1min |
| 6 | 2min |
| 7 | 3min |
| 8 | 4min |
| 9 | 5min |
| 10 | 6min |
| 11 | 7min |
| 12 | 8min |
| 13 | 9min |
| 14 | 10min |
| 15 | 20min |
| 16 | 30min |
| 17 | 1h |
| 18 | 2h |
发送延迟消息时,设置 delayLevel:
message.setDelayTimeLevel(16); // 30 分钟后投递
典型场景:订单创建后 30 分钟未支付自动取消。发送订单消息时设置 delayLevel=16(30min),Broker 延迟 30 分钟后投递,消费者检查订单状态,未支付则取消。
Kafka 不原生支持延迟消息,需要引入额外的时间轮组件(如 Quartz、Delay Queue),增加系统复杂度。
死信队列:消费失败的兜底
RocketMQ 原生支持死信队列(DLQ)。消息消费失败重试超过最大次数(默认 16 次)后,自动进入死信队列:
- 死信队列 Topic:
%DLQ% + ConsumerGroup - 消费者组专属:每个 ConsumerGroup 有独立的死信队列
- 人工处理:死信队列的消息需要人工排查和补偿
消费失败重试策略:
| 重试次数 | 延迟时间 |
|---|---|
| 第 1 次 | 10s |
| 第 2 次 | 30s |
| 第 3 次 | 1min |
| … | … |
| 第 16 次 | 2h |
超过 16 次后进入死信队列。
Kafka 不原生死信队列,需要自己实现:消费失败的消息写入专门的死信 Topic,人工处理。
长轮询:Push 模式的底层
RocketMQ 的 PushConsumer 不是真正的”推送”,而是长轮询(Long-Polling):
- Consumer 发起拉取请求
- Broker 若没有新消息,不立即返回空结果,而是将请求挂起(Hold)指定时间(如 15 秒)
- 在挂起期间,Producer 发来新消息,Broker 立刻唤醒请求返回数据
- 若超时无新消息,返回空结果
长轮询兼顾了 Push 的实时性和 Pull 的流控能力。RocketMQ 封装了消费线程池和负载均衡,对开发者屏蔽了细节,但也带来不透明的流控问题。
LitePullConsumer(RocketMQ 5.0+):完全由开发者控制拉取时机和数量,适合云原生场景下的精准流控。
常见误区
| 误区 | 真实情况 |
|---|---|
| RocketMQ 比 Kafka 吞吐高 | RocketMQ 单机万级到十万级 TPS,Kafka 单机十万级 TPS,吞吐是 Kafka 强项 |
| NameServer 像 ZooKeeper 做选主 | NameServer 只做路由注册,无状态,不做选主和协调 |
| 顺序消息是全局有序 | 只有 Queue 内有序,跨 Queue 无序 |
| 事务消息 = 分布式事务 | 事务消息保证本地事务 + 消息发送的一致性,不是跨服务的分布式事务 |
| 延迟消息任意时间 | 只支持 18 个固定延迟级别,不支持任意时间 |
项目判断
引入 RocketMQ 时,先回答三个问题:
- 业务场景:是否涉及订单、支付、库存这类”不能丢消息”的链路?
- 可靠性要求:是否需要事务消息、延迟消息、死信队列?
- Topic 数量:Topic 数量是否超过 50?
如果业务链路涉及交易、需要延迟消息、Topic 数量多 → RocketMQ 合适。如果只是日志采集、实时计算、吞吐优先 → Kafka 合适。
线上排查磁盘满错误(CODE: 14 DESC: service not available now, maybe disk full):
- 检查磁盘使用率,超过 90%(
diskMaxUsedSpaceRatio)Broker 拒绝写入 - 紧急清理过期 CommitLog 或迁移旧数据
- 调低阈值或设置自动清理时间(
deleteWhen=04) - 开启
transientStorePool(堆外内存缓存)减少磁盘依赖