Files
cs-note/hhs/MQ/07-主流MQ对比/24-RocketMQ.md
T
2026-05-24 20:51:06 +08:00

141 lines
8.0 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
---
tags: [MQ, RocketMQ, 消息队列, 分布式事务]
create time: 2026-05-24 19:52
---
# RocketMQ
## 概述
RocketMQ 是阿里巴巴开源的分布式消息中间件,最初为支撑淘宝双十一的万亿级消息流转而设计。它以高可用、高吞吐、低延迟著称,原生支持事务消息和延迟消息,在金融支付、电商交易、物流跟踪等场景中被广泛使用。本文从架构设计、核心特性、存储机制、集群部署到 5.0 新特性,系统梳理 RocketMQ 的知识体系。
## 正文
### 整体架构
RocketMQ 的架构由四个核心角色组成,各司其职:
```mermaid
graph TB
Producer["Producer 生产者"] -->|"发送消息"| NameServer
NameServer["NameServer 路由中心"] -->|"路由发现"| Consumer["Consumer 消费者"]
Producer -->|"写入消息"| Broker["Broker 消息代理"]
Consumer -->|"拉取消息"| Broker
Broker -->|"注册路由 + 心跳"| NameServer
style NameServer fill:#4A90D9,color:#fff
style Broker fill:#F5A623,color:#fff
style Producer fill:#6EC1E0,color:#fff
style Consumer fill:#6EC1E0,color:#fff
```
- **NameServer**:无状态的路由注册中心。Broker 启动时向所有 NameServer 注册路由信息(Topic 分布、Broker 地址等),Producer 和 Consumer 从 NameServer 拉取路由表。它不做任何消息中转,节点之间互不通信,水平扩展极为简单。
- **Broker**:消息存储和转发的核心。负责接收 Producer 发来的消息、持久化存储、响应 Consumer 的拉取请求。一个 Broker 可管理多个 Topic,每个 Topic 可设置多个 MessageQueue(类似 Kafka 的 Partition)。
- **Producer**:消息生产者,支持同步发送、异步发送和单向发送(Oneway)。
- **Consumer**:消息消费者,支持集群消费和广播消费两种模式。
> [!question] RocketMQ 的 NameServer 和 Kafka 的 ZooKeeper/KRaft 在职责上有什么本质区别?
> NameServer 是一个纯粹的路由注册表,只存储 Broker 的元数据,不做 Leader 选举、不做分区分配。而 ZooKeeper/KRaft 在 Kafka 中还承担 Controller 选举、Partition 副本分配、ISR 管理等职责。NameServer 的无状态设计让 RocketMQ 集群运维更简单,但路由信息的一致性是"最终一致"的(依赖 Broker 心跳),而非强一致。
### 核心特性
**事务消息**是 RocketMQ 最具辨识度的能力之一。它采用"半消息 + 本地事务 + 状态回查"三阶段机制:Producer 先发送一条半消息(Half Message)到 Broker,此时消息对 Consumer 不可见;Producer 执行本地事务后,根据结果提交(Commit)或回滚(Rollback);如果 Broker 长时间未收到确认,会主动回查 Producer 的本地事务状态。这套机制天然适合分布式事务的最终一致性场景。
**延迟消息**支持 18 个固定延迟级别(1s/5s/10s/30s/1m/2m...2h),底层通过定时任务 + 延迟队列实现。5.0 版本后新增任意时间精度的延迟消息,通过时间轮算法支撑更灵活的定时投递需求。
**消息过滤**分两层:Tag 过滤在 Broker 端完成,效率高但表达能力有限;SQL92 表达式过滤支持对消息属性做复杂条件判断(如 `amount > 100 AND region = 'CN'`),由 Consumer 端的 FilterServer 执行。
**消息回溯**允许 Consumer 按时间戳重置消费位点,重新消费历史消息,这在数据修复和问题排查中非常实用。
### 存储设计
RocketMQ 的存储模型采用三层结构:**CommitLog**(所有消息顺序追加写入一个大文件)+ **ConsumeQueue**(按 Topic-Queue 维度的逻辑索引,固定 20 字节/条)+ **IndexFile**(按消息 Key 的哈希索引,支持按 Key 查询消息)。
这种"先写大文件,再异步构建索引"的设计,将随机写转化为顺序写,最大化磁盘 IO 性能。ConsumeQueue 体积小,可常驻 PageCache,消费时几乎零磁盘 IO。
> [详细存储原理见 [[04-存储引擎/10-RocketMQ-CommitLog|RocketMQ CommitLog]]]
### 集群部署模式
RocketMQ 支持两种主从部署模式:
- **普通主从模式**:Master 接收读写请求,Slave 异步或同步复制数据。同步复制(SYNC_MASTER)保证数据不丢但延迟较高,异步复制(ASYNC_MASTER)延迟低但主节点宕机可能丢少量消息。
- **DLedger 模式**(推荐):基于 Raft 协议的自动选主模式。每个 Broker 组内 3 个节点,Leader 由 Raft 选举产生,日志通过 Raft 复制到多数节点后才提交。解决了普通主从模式下 Master 宕机需要人工介入的问题。
> [!question] 生产环境选同步复制还是异步复制?
> 看业务对数据可靠性的要求。金融交易场景选 SYNC_MASTER + SYNC_FLUSH(同步复制+同步刷盘),最大化数据安全;对延迟敏感但允许极少量消息丢失的场景选 ASYNC_MASTER + ASYNC_FLUSH。
### 消费模式
RocketMQ 的消费模式有两组维度:
| 维度 | 模式 | 说明 |
|------|------|------|
| 消费方式 | 集群消费(Clustering) | 同一 ConsumerGroup 内的实例分摊消息,每条消息只被消费一次 |
| | 广播消费(Broadcasting) | 每个 Consumer 实例都收到全量消息,适用于本地缓存刷新等场景 |
| 并发策略 | 并发消费 | 多线程同时消费,吞吐高但不保证顺序 |
| | 顺序消费 | 同一 MessageQueue 内的消息严格按顺序消费,通过 MessageListenerOrderly 实现 |
### RocketMQ 5.0 新特性
RocketMQ 5.0 带来了几项重要演进:
- **gRPC 协议**:替代自定义 RemotingCommand 协议,使用标准 gRPC 通信,多语言客户端支持更友好。
- **Pop 消费模式**:消费者不再需要维护 MessageQueue 的本地锁,由 Broker 端分配消息,更适应弹性伸缩和 Serverless 场景。
- **逻辑队列**:将物理 MessageQueue 抽象为逻辑队列,支持更灵活的队列映射和负载均衡策略。
### Go 代码示例
使用 `rocketmq-client-go` 发送和消费消息:
```go
// Producer: 同步发送消息
p, _ := rocketmq.NewProducer(
producer.WithNameServer([]string{"127.0.0.1:9876"}),
producer.WithGroupName("my-producer-group"),
)
p.Start()
defer p.Shutdown()
msg := rocketmq.NewMessage("OrderTopic", []byte(`{"orderId":"1001","amount":99.9}`))
msg.WithTag("order-create") // Tag 过滤标签
result, err := p.SendSync(context.Background(), msg)
fmt.Println("发送结果:", result.MessageID)
```
```go
// Consumer: 集群消费模式
c, _ := rocketmq.NewPushConsumer(
consumer.WithNameServer([]string{"127.0.0.1:9876"}),
consumer.WithGroupName("my-consumer-group"),
consumer.WithConsumeFromWhere(consumer.ConsumeFromLastOffset), // 从最新位点开始
)
c.Subscribe("OrderTopic", consumer.MessageSelector{
Type: consumer.TAG,
Expression: "order-create", // 只消费 Tag 为 order-create 的消息
}, func(ctx context.Context, msgs ...*ext.Message) (consumer.ConsumeResult, error) {
for _, msg := range msgs {
fmt.Println("收到消息:", string(msg.Body))
}
return consumer.ConsumeSuccess, nil
})
c.Start()
defer c.Shutdown()
// 阻塞等待消费
select {}
```
代码解析:Producer 通过 `SendSync` 同步发送消息到 `OrderTopic`,并通过 `WithTag` 设置 Tag。Consumer 通过 `Subscribe` 订阅同一 Topic,使用 Tag 表达式做过滤,消费成功返回 `ConsumeSuccess`。`ConsumeFromLastOffset` 表示新加入的消费者从最新消息开始消费,避免历史消息堆积。
## 关联笔记
- [[04-存储引擎/8-MQ-存储引擎设计|MQ 存储引擎设计]]
- [[04-存储引擎/10-RocketMQ-CommitLog|RocketMQ CommitLog]]
- [[06-高级特性/18-MQ-事务消息|MQ 事务消息]]
- [[06-高级特性/17-MQ-延迟消息与定时消息|MQ 延迟消息与定时消息]]
- [[06-高级特性/19-MQ-消息过滤与路由|MQ 消息过滤与路由]]
- [[05-可靠性保障/15-MQ-顺序性保障|MQ 顺序性保障]]
- [[12-架构与实战/44-MQ-高可用架构|MQ 高可用架构]]
- [[12-架构与实战/45-MQ-跨集群复制与容灾|MQ 跨集群复制与容灾]]
- [[07-主流MQ对比/27-MQ-选型对比|MQ 选型对比]]