--- 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 选型对比]]