RocketMQNotes

第 01 章:认识 RocketMQ

zjc 于 2026-01-01 发布

这是《RocketMQ 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 RocketMQ 是阿里巴巴开源的分布式消息中间件,后来成为 Apache 顶级项目。它以业务消息见长,广泛用于订单事件、交易解耦、异步通知、削峰填谷、事务消息、延迟消息和流式集成。

如果你已经读过 Kafka,可以把 RocketMQ 理解为另一种消息系统取舍:Kafka 更偏向高吞吐日志流和流处理生态,RocketMQ 更强调业务消息语义,例如事务消息、延迟消息、重试队列、死信队列和消息轨迹。

1.1 消息系统解决什么问题

假设订单服务创建订单后要做六件事:

创建订单
  -> 扣库存
  -> 发优惠券
  -> 发短信
  -> 更新积分
  -> 同步搜索索引
  -> 风控审计

同步调用的弊端:

  1. 任何一个下游失败都会影响下单;
  2. 下游耗时累加,接口延迟变高;
  3. 流量高峰时容易雪崩;
  4. 新增下游要修改订单服务;
  5. 失败重试和状态恢复难以统一治理。

引入消息后:

订单服务
  -> 保存订单
  -> 发送 OrderCreated 事件

下游消费者
  -> 库存服务
  -> 营销服务
  -> 通知服务
  -> 搜索同步服务
  -> 风控服务

消息系统的价值:

能力 说明
解耦 生产者不需要知道所有下游
异步 主链路只做核心事务
削峰 突发流量先进入队列
重试 消费失败自动重试
死信 异常消息可隔离治理
广播 一条事件被多个消费组消费
顺序 同一 key 的消息可按顺序消费
审计 事件可作为系统事实流

1.2 RocketMQ 的核心角色

RocketMQ 4.x 主要由 NameServer、Broker、Producer、Consumer 组成。

flowchart LR
    P[Producer] -->|发送消息| B[Broker]
    C[Consumer] -->|拉取消息| B
    NS[NameServer] -.路由发现.- P
    NS -.路由发现.- C
    NS -.心跳注册.- B

NameServer

NameServer 是轻量级注册中心,负责保存 Topic 路由信息。

特点:

Broker

Broker 是消息存储和服务节点,负责:

Producer

生产者负责发送消息。常见发送方式:

模式 说明
同步发送 等待 Broker 确认,可靠性高
异步发送 回调处理结果,延迟低
单向发送 不等待结果,允许丢失

Consumer

消费者从 Broker 拉取消息并执行业务处理。消费模式:

模式 说明
集群消费 一个消费组分摊队列,适合业务事件
广播消费 每个消费者都消费全量,适合本地缓存刷新

1.3 Topic、Queue 与消费组

RocketMQ 的逻辑模型:

Topic
  +-- MessageQueue 0
  +-- MessageQueue 1
  +-- MessageQueue 2

一个 Topic 有多个 MessageQueue。生产者选择队列写入,消费者组中的实例分摊队列。

消费位点

每个消费组对每个队列维护一个消费位点 offset。位点表示:

这个消费组已经消费到哪里

重要原则:

  1. 业务处理成功后再提交位点;
  2. 消费失败会进入重试;
  3. 位点不能随意重置;
  4. 消费组隔离,不同组互不影响;
  5. 位点积压是最重要的监控指标。

1.4 RocketMQ 5.x 架构变化

RocketMQ 5.x 引入了 Proxy 和新的客户端架构。

Client
  -> Proxy
     -> Broker

Proxy 的价值:

  1. 客户端接入层标准化;
  2. 支持 gRPC 协议;
  3. 简化多语言客户端;
  4. 支持云原生部署形态;
  5. 将接入逻辑与存储逻辑分离。

5.x 仍保留对 4.x Java 客户端的兼容,但新项目应评估:

1.5 RocketMQ 的消息类型

普通消息

无顺序要求、无定时要求的一般业务事件。

顺序消息

同一 Sharding Key 的消息进入同一队列,消费者对队列串行处理。

典型场景:

延迟消息

消息发送后延迟一段时间才可消费。

典型场景:

事务消息

解决“本地事务和消息发送”的一致性问题。

典型流程:

发送 half 消息
  -> 执行本地事务
  -> 提交或回滚消息
  -> Broker 回查未知状态

事务消息不是分布式事务的万能方案,它保证的是“本地事务成功后消息最终可见”。

1.6 存储模型概览

RocketMQ 的核心存储文件:

文件 作用
CommitLog 所有 Topic 消息顺序写入的主文件
ConsumeQueue 每个队列的逻辑索引
IndexFile 消息 Key 或时间查询索引
Checkpoint 刷盘和恢复位点
abort 异常停机标记

写入路径:

Producer 消息
  -> Broker 接收
  -> 写 CommitLog
  -> 构建 ConsumeQueue
  -> 刷盘
  -> 主从复制

所有消息追加到 CommitLog,可以让磁盘写入保持顺序 IO,这是 RocketMQ 高吞吐的重要基础。

1.7 和 Kafka 的对比

维度 Kafka RocketMQ
定位 日志流、流处理生态 业务消息中间件
存储 Partition 分文件日志 CommitLog + ConsumeQueue
事务 流事务语义 本地事务消息
延迟消息 需自行实现 原生支持
重试队列 消费者自行处理 原生支持
死信队列 需自行实现 原生支持
消息轨迹 依赖生态 原生能力
流处理 Kafka Streams / Flink 生态强 相对偏业务集成

选择建议:

  1. 日志采集、流计算、高吞吐管道优先考虑 Kafka;
  2. 订单、交易、营销等业务事件优先考虑 RocketMQ;
  3. 两者也可以共存,不要强行用一个系统覆盖所有场景;
  4. 迁移时要评估语义、运维、客户端生态和团队经验。

1.8 生产使用的基本要求

一个生产级 RocketMQ 集群至少要有:

常见红线:

  1. 消息发送失败不能只打印日志;
  2. 消费者不能吞掉异常;
  3. 不能随意重置消费位点;
  4. 不能让死信队列无限增长;
  5. Topic 扩队列要评估顺序性影响;
  6. 生产者不能无限重试造成雪崩;
  7. 消费逻辑不能无超时;
  8. 关键业务不能使用单向发送。

本章小结

RocketMQ 是面向业务事件的分布式消息系统,核心角色包括 NameServer、Broker、Producer 和 Consumer。它用 CommitLog 顺序写保证吞吐,用 ConsumeQueue维护消费索引,并提供事务消息、延迟消息、顺序消息、重试队列和死信队列等业务语义。理解和 Kafka 的差异,才能在架构中正确选择。

思考题

  1. RocketMQ 的 NameServer 和 Kafka 的元数据机制有什么差异?
  2. CommitLog 和 ConsumeQueue 分别解决什么问题?
  3. 什么场景必须使用顺序消息?顺序消息会牺牲什么?
  4. 事务消息解决的是哪个一致性问题?它不是什么?
  5. 如何设计消息消费幂等?