这是《RocketMQ 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 RocketMQ 是阿里巴巴开源的分布式消息中间件,后来成为 Apache 顶级项目。它以业务消息见长,广泛用于订单事件、交易解耦、异步通知、削峰填谷、事务消息、延迟消息和流式集成。
如果你已经读过 Kafka,可以把 RocketMQ 理解为另一种消息系统取舍:Kafka 更偏向高吞吐日志流和流处理生态,RocketMQ 更强调业务消息语义,例如事务消息、延迟消息、重试队列、死信队列和消息轨迹。
1.1 消息系统解决什么问题
假设订单服务创建订单后要做六件事:
创建订单
-> 扣库存
-> 发优惠券
-> 发短信
-> 更新积分
-> 同步搜索索引
-> 风控审计
同步调用的弊端:
- 任何一个下游失败都会影响下单;
- 下游耗时累加,接口延迟变高;
- 流量高峰时容易雪崩;
- 新增下游要修改订单服务;
- 失败重试和状态恢复难以统一治理。
引入消息后:
订单服务
-> 保存订单
-> 发送 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
Broker 是消息存储和服务节点,负责:
- 接收生产者消息;
- 持久化 CommitLog;
- 维护 ConsumeQueue;
- 处理消费者拉取;
- 管理消费位点;
- 处理重试和死信;
- 主从同步;
- 管理延迟消息和事务回查。
Producer
生产者负责发送消息。常见发送方式:
| 模式 | 说明 |
|---|---|
| 同步发送 | 等待 Broker 确认,可靠性高 |
| 异步发送 | 回调处理结果,延迟低 |
| 单向发送 | 不等待结果,允许丢失 |
Consumer
消费者从 Broker 拉取消息并执行业务处理。消费模式:
| 模式 | 说明 |
|---|---|
| 集群消费 | 一个消费组分摊队列,适合业务事件 |
| 广播消费 | 每个消费者都消费全量,适合本地缓存刷新 |
1.3 Topic、Queue 与消费组
RocketMQ 的逻辑模型:
Topic
+-- MessageQueue 0
+-- MessageQueue 1
+-- MessageQueue 2
一个 Topic 有多个 MessageQueue。生产者选择队列写入,消费者组中的实例分摊队列。
消费位点
每个消费组对每个队列维护一个消费位点 offset。位点表示:
这个消费组已经消费到哪里
重要原则:
- 业务处理成功后再提交位点;
- 消费失败会进入重试;
- 位点不能随意重置;
- 消费组隔离,不同组互不影响;
- 位点积压是最重要的监控指标。
1.4 RocketMQ 5.x 架构变化
RocketMQ 5.x 引入了 Proxy 和新的客户端架构。
Client
-> Proxy
-> Broker
Proxy 的价值:
- 客户端接入层标准化;
- 支持 gRPC 协议;
- 简化多语言客户端;
- 支持云原生部署形态;
- 将接入逻辑与存储逻辑分离。
5.x 仍保留对 4.x Java 客户端的兼容,但新项目应评估:
- 客户端版本;
- Proxy Local / Cluster 模式;
- Controller 模式;
- 存算分离能力;
- 与现有运维工具的兼容性。
1.5 RocketMQ 的消息类型
普通消息
无顺序要求、无定时要求的一般业务事件。
顺序消息
同一 Sharding Key 的消息进入同一队列,消费者对队列串行处理。
典型场景:
- 同一订单的状态流;
- 同一用户的积分流水;
- 同一设备的指令顺序。
延迟消息
消息发送后延迟一段时间才可消费。
典型场景:
- 订单 30 分钟未支付自动取消;
- 延迟检查任务;
- 定时提醒。
事务消息
解决“本地事务和消息发送”的一致性问题。
典型流程:
发送 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 生态强 | 相对偏业务集成 |
选择建议:
- 日志采集、流计算、高吞吐管道优先考虑 Kafka;
- 订单、交易、营销等业务事件优先考虑 RocketMQ;
- 两者也可以共存,不要强行用一个系统覆盖所有场景;
- 迁移时要评估语义、运维、客户端生态和团队经验。
1.8 生产使用的基本要求
一个生产级 RocketMQ 集群至少要有:
- 多个 NameServer;
- Broker 主从或多副本;
- 监控消息积压;
- 监控发送失败率;
- 监控消费失败和死信;
- 消息发送幂等;
- 消费幂等;
- 死信人工处理流程;
- 容量规划和压测;
- 故障降级预案。
常见红线:
- 消息发送失败不能只打印日志;
- 消费者不能吞掉异常;
- 不能随意重置消费位点;
- 不能让死信队列无限增长;
- Topic 扩队列要评估顺序性影响;
- 生产者不能无限重试造成雪崩;
- 消费逻辑不能无超时;
- 关键业务不能使用单向发送。
本章小结
RocketMQ 是面向业务事件的分布式消息系统,核心角色包括 NameServer、Broker、Producer 和 Consumer。它用 CommitLog 顺序写保证吞吐,用 ConsumeQueue维护消费索引,并提供事务消息、延迟消息、顺序消息、重试队列和死信队列等业务语义。理解和 Kafka 的差异,才能在架构中正确选择。
思考题
- RocketMQ 的 NameServer 和 Kafka 的元数据机制有什么差异?
- CommitLog 和 ConsumeQueue 分别解决什么问题?
- 什么场景必须使用顺序消息?顺序消息会牺牲什么?
- 事务消息解决的是哪个一致性问题?它不是什么?
- 如何设计消息消费幂等?