Skip to content

RocketMQ 详解:架构、事务消息与延时消息

本文是消息队列系统学习系列的 L2 核心篇。前置:28. kafka-producer-consumer-internals29. kafka-reliability-high-performance。 学完可以配合面试题食用:13-rocketmq-architecture-nameserver-broker-producer-consumer14-rocketmq-transaction-message16-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 -->|拉消息| B3

NameServer 做了什么? 每个 NameServer 独立维护全量路由表,无状态。Broker 启动后向所有 NameServer 注册 Topic 路由(Broker 地址 + 队列数),每 30 秒心跳续命。Producer/Consumer 启动时从 NameServer 拉路由,缓存在本地。NameServer 挂一个不影响集群,因为客户端走的是本地缓存,且会尝试其他 NameServer。

和 Kafka 有什么不同?

维度KafkaRocketMQ
元数据服务ZooKeeper / KRaftNameServer(无状态)
存储单元Partition → SegmentConsumeQueue + 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 秒延时怎么办? 两种方案:

  1. 取 10s(level 4),应用层做 3 秒 Timer 补偿
  2. 自己用时间轮实现:在内存里维护一个延时队列,到期后发 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 放在一起对比,给出具体场景的选型结论。

参考

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。
粤ICP备2026104257号-1