消息队列02:RocketMQ

从 NameServer 路由、Broker 存储模型、事务消息两阶段提交到顺序消息实现,拆解 RocketMQ 在业务消息场景下的设计逻辑和可靠性保证。

字数 2851 阅读时长 ≈ 9 分钟 2026-7-13 2026-7-13
消息队列02:RocketMQ

RocketMQ 不是为”日志管道”设计的,而是为金融级业务消息设计的。它的核心场景是订单、支付、库存同步这类”不能丢消息、要保证顺序、需要事务联动”的链路。

这篇文章围绕四个核心机制:NameServer、Broker 存储模型、事务消息、顺序消息,把 RocketMQ 与 Kafka 的差异和业务适配讲清楚。

为什么业务消息选 RocketMQ

先定位清楚:RocketMQ 的核心场景是金融级可靠事务消息,追求消息零丢失和严格顺序,适合订单、交易、库存等核心业务链路。

与 Kafka 的核心差异:

维度RocketMQKafka
存储模型集中式 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/OPartition 过多时退化为随机写
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 最大的差异化能力:保证本地事务和消息发送要么同时成功,要么同时失败

事务消息流程

事务消息采用两阶段提交 + 事务状态回查

  1. 发送 Half 消息:Producer 发送半消息(Half Message)到 Broker,消息对消费者不可见,Broker 返回发送成功
  2. 执行本地事务:Producer 执行本地业务逻辑(如写入数据库)
  3. 提交或回滚
    • 本地事务成功 → Producer 发送 Commit 消息,Broker 将 Half 消息变为可消费
    • 本地事务失败 → Producer 发送 Rollback 消息,Broker 删除 Half 消息
  4. 事务回查:如果 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 个延迟级别:

级别延迟时间
11s
25s
310s
430s
51min
62min
73min
84min
95min
106min
117min
128min
139min
1410min
1520min
1630min
171h
182h

发送延迟消息时,设置 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)

  1. Consumer 发起拉取请求
  2. Broker 若没有新消息,不立即返回空结果,而是将请求挂起(Hold)指定时间(如 15 秒)
  3. 在挂起期间,Producer 发来新消息,Broker 立刻唤醒请求返回数据
  4. 若超时无新消息,返回空结果

长轮询兼顾了 Push 的实时性和 Pull 的流控能力。RocketMQ 封装了消费线程池和负载均衡,对开发者屏蔽了细节,但也带来不透明的流控问题。

LitePullConsumer(RocketMQ 5.0+):完全由开发者控制拉取时机和数量,适合云原生场景下的精准流控。

常见误区

误区真实情况
RocketMQ 比 Kafka 吞吐高RocketMQ 单机万级到十万级 TPS,Kafka 单机十万级 TPS,吞吐是 Kafka 强项
NameServer 像 ZooKeeper 做选主NameServer 只做路由注册,无状态,不做选主和协调
顺序消息是全局有序只有 Queue 内有序,跨 Queue 无序
事务消息 = 分布式事务事务消息保证本地事务 + 消息发送的一致性,不是跨服务的分布式事务
延迟消息任意时间只支持 18 个固定延迟级别,不支持任意时间

项目判断

引入 RocketMQ 时,先回答三个问题:

  1. 业务场景:是否涉及订单、支付、库存这类”不能丢消息”的链路?
  2. 可靠性要求:是否需要事务消息、延迟消息、死信队列?
  3. Topic 数量:Topic 数量是否超过 50?

如果业务链路涉及交易、需要延迟消息、Topic 数量多 → RocketMQ 合适。如果只是日志采集、实时计算、吞吐优先 → Kafka 合适。

线上排查磁盘满错误(CODE: 14 DESC: service not available now, maybe disk full):

  • 检查磁盘使用率,超过 90%(diskMaxUsedSpaceRatio)Broker 拒绝写入
  • 紧急清理过期 CommitLog 或迁移旧数据
  • 调低阈值或设置自动清理时间(deleteWhen=04
  • 开启 transientStorePool(堆外内存缓存)减少磁盘依赖