1
0

RocketMQ 笔记

2026-04-12
2026-07-23

一、RocketMQ 是什么,为什么微服务需要消息队列

1.1RocketMQ简介

RocketMQ 是一个分布式消息与事件流平台,核心作用是在不同应用、服务或模块之间传递消息。它不会替代 Dubbo、OpenFeign 或 HTTP,而是解决它们不擅长的问题:当调用方不需要立即得到处理结果,或者业务希望削峰、解耦、异步处理时,可以不直接同步调用下游服务,而是先把消息发送到 RocketMQ,由下游服务在合适的时间消费。

同步调用的典型链路如下:

订单服务
    ↓ 同步调用
库存服务
    ↓ 同步调用
积分服务
    ↓ 同步调用
通知服务

如果每个调用耗时 100ms,链路越长,用户等待越久,而且任何一个下游服务故障都可能导致整个请求失败。订单服务还必须知道库存、积分、通知等服务的地址和接口,服务之间耦合较强。

使用 RocketMQ 后,可以改成:

用户提交订单
      ↓
订单服务保存订单
      ↓
发送 OrderCreatedEvent
      ↓
RocketMQ Broker 持久化消息
      ↓
 ┌────────────┬────────────┬────────────┐
库存服务      积分服务      通知服务
扣减库存      增加积分      发送短信

订单服务只负责发布“订单已创建”这一事实,并不需要同步等待所有下游任务完成。库存、积分、通知可以分别订阅同一个 Topic,各自消费和处理。RocketMQ 官方领域模型把一条消息的生命周期分为生产、存储和消费三个阶段:Producer 产生消息,Broker 将消息保存到 Topic,Consumer 订阅并消费消息。

消息队列最核心的三个作用是 异步、解耦和削峰

异步是把不需要立即返回结果的操作从当前请求线程中移走。例如用户上传文件后,文件元数据需要立即保存,但缩略图生成、内容审核、向量化、消息通知不一定要在同一次 HTTP 请求中完成。主流程只要保存文件并发送消息即可,后续任务由消费者异步完成。

解耦是让生产者不直接依赖消费者。订单服务不必逐个调用积分、短信和数据分析服务,只需要发送统一的订单事件。未来增加风控服务时,只要让风控服务订阅 Topic,通常不需要修改订单服务。

削峰是用 Broker 的持久化能力暂时接住流量。秒杀活动瞬间产生十万条请求,数据库可能每秒只能稳定处理一万次写入。可以先把请求转换成消息写入 RocketMQ,消费者再按照数据库能够承受的速度处理。消息队列不是把流量消灭,而是把短时间内的流量峰值摊平到更长时间。

除此之外,RocketMQ 还适合事件驱动、最终一致性、延时任务、顺序处理、事务消息、日志采集和数据同步等场景。它不适合所有调用:用户登录校验、余额查询、读取文件详情等必须立即得到结果的操作,仍然更适合使用 HTTP、Dubbo 或数据库查询。

1.2为什么要使用MQ

1,要做到系统解耦,当新的模块进来时,可以做到代码改动最小;  能够解耦 2,设置流程缓冲池,可以让后端系统按自身吞吐能力进行消费,不被冲垮; 能够削峰,限流 3,强弱依赖梳理能把非关键调用链路的操作异步化并提升整体系统的吞吐能力;能够异步

Mq的作用  削峰限流 异步 解耦合

1.2.1 定义

  • MQ是一个消息中间件(如缓存中间件有redis、memcache;数据库中间件有mycat、canal)
  • 利用高效可靠的消息传递机制进行与平台无关(跨平台)的数据交流,并基于数据通信来进行分布式系统的集成。
  • 通过提供消息传递和消息排队模型在分布式环境下提供应用解耦,弹性伸缩,冗余存储,流量削峰,异步通信,数据同步等

大致流程

  • 发送者把消息发给消息服务器 MQ,消息服务器把消息存放在若干队列/主题中,在合适的时候,消息服务器会把消息转发给接受者。在这个过程中,发送和接受是异步的,也就是发送无需等待,发送者和接受者的生命周期也没有必然关系在发布pub/订阅sub模式下
  • 在发布pub/订阅sub模式下,也可以完成一对多的通信,即让一个消息有多个接受者(例如微信订阅号) ![](file:////tmp/wps-madm/ksohtml/wpsNnbEfo.jpg)

1.2.2 MQ作用

为什么要用消息队列来存储消息,生产者为什么不直接发给消费者?

[!note] 作用一:异步

消息发送者可以发送一个消息而无需等待响应。消息发送者把消息发送到一条虚拟的通道(主题或队列)上;

消息接收者则订阅或监听该通道。一条信息可能最终转发给一个或多个消息接收者,这些接收者都无需对消息发送者做出回应。整个过程都是异步的。

案例:

也就是说,一个系统和另一个系统间进行通信的时候,假如系统A希望发送一个消息给系统B,让它去处理,但是系统A不关注系统B到底怎么处理或者有没有处理好,所以系统A把消息发送给MQ,然后就不管这条消息的“死活” 了,接着系统B从MQ里面消费出来处理即可。至于怎么处理,是否处理完毕,什么时候处理,都是系统B的事,与系统A无关。

![](file:////tmp/wps-madm/ksohtml/wpsfmvozq.jpg)

这样的一种通信方式,就是所谓的“异步”通信方式,对于系统A来说,只要把消息发给MQ,然后系统B就会异步处去进行处理了,系统A不需要“同步”的等待系统B处理完。这样的好处是什么呢?解耦

[!note] 作用二:削峰限流

消费者的业务执行能力是有限的,例如一个业务执行需要花费3秒。假如大量请求同时过来,tomcat线程池的线程被占用完了,后续的请求会报错(503服务不可用)。如果请求过来之后,直接放到消息队列就返回,这样就可以大大降低服务不可用的概率,后续消费者再去消息队列里面慢慢获取业务信息去执行即可;

[!note] 作用三:解耦

取消应用之间、模块之间的强耦合关系,让程序更加健壮、高可用。如果模块都直接耦合在一块,一个模块出问题,很容易影响到其他模块。通过使用A模块->消息队列->B模块模式将A、B模块解耦

【应用系统解耦的好处】

  • 发送者和接收者不必了解对方,只需要确认消息
  • 发送者和接收者不必同时在线

现实中的业务 ![](file:////tmp/wps-madm/ksohtml/wpsQclBUo.png)

1.2.3 各个MQ产品的比较

MQ主要关注两个性能:

  • 吞吐量:单位时间内可以处理多少条消息(用消息大小来描述更准确,因为消息内存越小,数量会越大)

  • 时效性:生产者发消息,MQ多久才收到消息

【对比】

特性 ActiveMQ RabbitMQ Rocket MQ kafka
开发语言 java erlang java scala
单机吞吐量 万级 万级 10万级 10万级
时效性 ms级 us级 ms级 ms级以内
可用性 高(主从架构) 高(主从架构) 非常高(分布式架构) 非常高(分布式架构)
功能特性 成熟的产品,在很多公司得到应用;有较多的文档;各种协议支持较好 基于Erlang开发,所以并发能力很强,性能极其好,延时很低;管理界面较丰富 MQ功能比较完备,扩展性佳 只支持主要的MQ功能,像一些消息查询,消息回溯等功能没有提供,毕竟是为大数据准备的,在大数据领域应用广。

【总结】

  • activeMQ:使用java实现(jms 协议),性能一般,出现早,功能单一,吞吐量低
  • rabbitmq:使用erlang实现(amqp 协议),性能好,功能丰富,吞吐量一般
  • rocketmq:使用java实现,性能好,功能最丰富,吞吐量高
  • kafka:使用scala实现,吞吐量最大,功能单一(专注读、写),主要用于大数据领域(数据又多又大)

1.3 RocketMQ重要概念

RocketMQ是阿里巴巴2016年MQ中间件,使用Java语言开发,是一款开源的分布式消息系统(可以做集群),基于高可用分布式集群技术,提供低延时的、高可靠的消息发布与订阅服务。同时,广泛应用于多个领域,包括异步通信解耦、企业解决方案、金融支付、电信、电子商务、快递物流、广告营销、社交、即时通信、移动应用、手游、视频、物联网、车联网等。

具有以下特点:

  1. 能够严格保证消息按照顺序消费
  2. 提供丰富的消息拉取模式
  3. 高效的订阅者水平扩展能力
  4. 实时的消息订阅机制
  5. 亿级消息堆积能力

[!note] 相关重要概念

  • Producer(生产者):消息的发送者;举例:发件人
  • Consumer(消费者):消息的接收者;举例:收件人
  • Broker:暂存和传输消息的通道;举例:快递公司
  • NameServer:管理Broker,相当于broker的注册中心,保留了broker的信息;举例:各个快递公司的管理机构
  • Queue(队列):消息真实存放的位置,一个Broker中可以有多个队列
  • Topic(主题):消息的分类(虚拟的结构,用来区分不同类型的消息)
  • ProducerGroup(生产者组)
  • ConsumerGroup(消费者组):多个消费者组可以同时消费一个主题的消息 消息发送的流程是,Producer询问NameServer,NameServer分配一个broker 然后Consumer也要询问NameServer,得到一个具体的broker,然后消费消息 [单机版本结构] ![](file:////tmp/wps-madm/ksohtml/wpse4fqUn.jpg)

1.3.1 RocketMQ为何高可用

如何设计高可用消息队列? 只部署一个 Broker 并不能称为高可用。一旦该 Broker 宕机,生产者无法发送消息,消费者也无法继续拉取消息,整个消息链路都会中断。

1 使用broker集群实现写的高可用

多个 Master Broker 组成 Broker 集群,共同承载一个 Topic 的多个 MessageQueue:

Topic:OrderTopic(共 6 个 Queue)
              │
          ┌───┴───┐
          ▼       ▼
       Broker A  Broker B
       Queue 0   Queue 3
       Queue 1   Queue 4
       Queue 2   Queue 5
┌─────────────────────────────────────────────────────────────┐
│                    Broker 集群(Cluster)                     │
│                                                             │
│   ┌─────────────────────┐      ┌─────────────────────┐      │
│   │    Broker A 主从     │      │    Broker B 主从    │      │
│   │  ┌─────┐ ┌─────┐    │      │  ┌─────┐ ┌─────┐    │      │
│   │  │Master│ │Slave│   │      │  │Master│ │Slave│   │      │
│   │  │(读写)│ │(同步)│    │      │  │(读写)│ │(同步)│   │      │
│   │  └──┬──┘ └──┬──┘    │      │  └──┬──┘ └──┬──┘    │      │
│   │     └────同步────┘   │      │     └────同步────┘  │       │
│   │存 Topic 的 Queue 0-2 │      │存 Topic 的 Queue 3-5 │      │
│   └─────────────────────┘      └─────────────────────┘      │
│                                                             │
│             ←── 集群:分散负载、水平扩展 ──→                     │
│             ←── 主从:数据备份、故障切换 ──→                    │
└─────────────────────────────────────────────────────────────┘

需要注意,多 Master 并不等于任意一个 Master 宕机后所有写请求都不受影响。 不同 Master 保存的是不同 MessageQueue。当 Broker A 宕机后:

  • Broker A 上的 Queue 0~2 暂时不可写、不可直接访问;
  • Producer 可以在重试时选择 Broker B 上仍然可用的 Queue;
  • 因此 Topic 整体通常仍可继续接收消息;
  • 但 Broker A 上原有 Queue 的消息需要等待节点恢复,或者由其 Slave 提供读取;
  • 如果业务指定固定 Queue、使用严格顺序消息或自定义 Queue 选择器,则不一定能够自动切换到其他 Broker。

因此,多 Master 主要解决的是:

单个 Broker 故障时,Topic 仍有其他 Broker 和 Queue 可以承载新消息。

它提高了 Topic 级别的写可用性和系统吞吐量,但不能恢复故障 Broker 上原有 Queue 的写能力。

多个 Master 节点分别存储 Topic 的部分 Queue,能够分散请求、提高吞吐量。某个 Master 宕机后,其他 Master 上的 Queue 仍然可以继续收发消息,但故障节点上的 Queue 暂时不可写;是否能够自动重试到其他 Queue,还取决于发送方式、重试机制以及是否要求严格顺序。

2 使用主从复制提高数据可靠性和读可用性
    Master Broker A          Slave Broker A
    (主,读写)               (从,同步复制)
         │                        │
         └──────── 同步复制 ──────┘
  • Master:负责接收 Producer 写入,也负责处理 Consumer 的读取请求。
  • Slave:从 Master 复制 CommitLog 数据,不能接收 Producer 写入,但可以在特定条件下承担读请求。
  • Master 和 Slave 使用相同的 brokerName,通过不同的 brokerId 区分:brokerId=0 表示 Master,非 0 表示 Slave。

RocketMQ 客户端可以同时感知 Master 和 Slave。正常情况下通常优先从 Master 拉取消息;当 Master 压力较大或不可用时,Broker 可以建议 Consumer 从 Slave 拉取数据,因此传统主从架构能够提高消费端的读可用性

但是,传统主从架构下 Slave 不能接收写入。Master 宕机后,即使 Consumer 还可以从 Slave 读取已经同步的消息,Producer 也不能把新消息写入该 Slave。要恢复写入,必须人工切换主从,或者使用 DLedger、Controller 等自动选主机制。

主从复制分为两种:

复制方式 特点 适用场景
同步复制 SYNC_MASTER Master 将消息传输到 Slave,并收到 Slave 确认后再向 Producer 返回成功 数据可靠性要求较高,可以接受更高延迟
异步复制 ASYNC_MASTER Master 写入本机后即可返回,Slave 在后台异步复制 追求吞吐量,能够容忍 Master 突然故障时极少量未复制消息丢失

**RocketMQ 常见默认配置是:

brokerRole=ASYNC_MASTER
flushDiskType=ASYNC_FLUSH

即:

异步主从复制 + 异步刷盘

这种配置性能较高,但 Broker 所在机器突然断电或磁盘损坏时,理论上可能丢失少量尚未刷盘或尚未同步到 Slave 的消息。

对于可靠性要求更高的业务,可以使用:

brokerRole=SYNC_MASTER
flushDiskType=SYNC_FLUSH

即同步复制加同步刷盘,但写入延迟和吞吐量都会受到影响。

Slave 的读取能力可以概括为:

场景 Slave 是否能够读取
Master 正常 可以根据 Broker 和客户端负载策略,从 Slave 拉取部分消息
Master 宕机 Consumer 可以继续读取 Slave 中已经同步的数据
Producer 写入 不可以,传统 Slave 始终不能接受生产写入
主从复制延迟较大 从 Slave 读取时可能暂时读不到 Master 上最新写入、但尚未完成复制的消息
因此,主从复制主要解决两个问题:
数据副本备份;
Master 不可用时,Consumer 仍可读取 Slave 中已经复制的消息。

它并不等于传统 Slave 可以自动接管写流量。

3 Broker 集群与主从复制组合使用

实际生产环境通常将多 Master 集群与主从复制结合使用:

部署方式 说明 适用场景
1. 单 Master 模式 只有一个 Broker,没有 Slave。部署最简单,但存在明显单点故障;Broker 宕机后,消息生产和消费均不可用,且没有副本备份 本地开发、功能测试,不建议用于生产
2. 多 Master 模式,无 Slave 多台独立 Master,Topic 的 Queue 分散在不同 Broker。没有主从复制开销,吞吐量较高;单个 Master 宕机后,其他 Master 的 Queue 仍可收发消息,但故障节点上的 Queue 暂时不可用,磁盘损坏还可能造成消息永久丢失 日志采集、埋点等允许部分 Queue 暂时不可用的非核心业务
3. 多 Master、多 Slave,传统主从 每个 Master 配置一个或多个 Slave。Master 宕机后,Consumer 可以从 Slave 读取已经同步的消息,但 Slave 默认不会自动升级为 Master,新消息仍不能写入该主从组。异步复制性能较高,但可能丢失极少量未同步消息;同步复制可靠性更高,但延迟和吞吐量有所下降 一般生产业务,如订单通知、异步任务、业务解耦等
4. DLedger 自动切换模式 一组 Broker 通常至少部署 3 个节点,使用基于 Raft 的 DLedger 进行日志复制和 Leader 选举。Leader 故障后,可以重新选举 Leader 并恢复写服务,不需要人工修改 brokerId。选举与恢复需要一定时间,不能严格表述为固定 10 秒或完全无缝 RocketMQ 4.x 中需要自动故障切换的高可靠部署
5. RocketMQ 5.x Controller 模式 在传统 Broker 副本架构上增加 Controller,由 Controller 管理副本状态和主节点选举。Master 故障后,可自动从同步状态较好的副本中选出新 Master。Controller 可以独立部署,也可以嵌入部分 NameServer RocketMQ 5.x 核心生产系统,需要自动主从切换和较高可用性的业务

需要特别注意:

传统主从:Master 故障后,Consumer 可以切换读取 Slave,
          但 Slave 不会自动成为 Master,写服务不能自动恢复。

DLedger / Controller:发生故障后可以自动选主,
                     从而恢复该 Broker 副本组的读写能力。

生产环境典型部署:

                     NameServer 集群
                            │
                  ┌─────────┴─────────┐
                  ▼                   ▼
          ┌─────────────┐     ┌─────────────┐
          │ Broker A 组  │     │ Broker B 组  │
          │ Master       │     │ Master       │
          │    +         │     │    +         │
          │ Slave        │     │ Slave        │
          └─────────────┘     └─────────────┘
                 ▲                   ▲
                 │                   │
          主从复制与数据备份    主从复制与数据备份

          ←──── 多 Broker:分散 Queue 和请求负载 ────→
          ←──── 主从副本:提高数据与读取可用性 ─────→

RocketMQ 的高可用不是由某一个机制单独实现的,而是多个层次共同作用:

NameServer 集群:避免路由服务出现单点故障;
多 Master 集群:分散 Topic Queue,提高吞吐量和 Topic 整体可用性;
Master-Slave:保存数据副本,提高数据可靠性和消费端读可用性;
同步刷盘:降低机器断电导致的数据丢失风险;
同步复制:降低 Master 损坏导致的数据丢失风险;
DLedger / Controller:在 Master 故障后自动选主,恢复写服务;
Producer 重试:当前 Broker 或 Queue 发送失败时尝试其他可用 Queue。

因此,一个较完整的高可用 RocketMQ 架构通常需要:

NameServer 集群
       +
多 Broker 副本组
       +
主从数据复制
       +
合理的刷盘策略
       +
自动主从切换机制
       +
生产者重试和业务幂等

其中,Broker 集群解决的是负载分散和 Topic 整体可用性,主从复制解决的是数据副本与消费读取可用性,DLedger 或 Controller 解决的是Master 故障后的自动选主和写服务恢复

1.3.2Broker内部结构

一个Broker里面可以存放多个Topic,一个Topic里面可以存放多个队列

1.3.3生产和消费理解

1.3.3消费模式

MQ的消费模式可以大致分为两种

  • 推 Push :服务端【MQ】主动推送消息给客户端,优点是及时性较好,但如果客户端没有做好流控,一旦服务端推送大量消息到客户端时,就会导致客户端消息堆积甚至崩溃(客户端压力大
  • 拉 Pull :客户端需要主动到服务端【MQ】拉取数据,优点是客户端可以依据自己的消费能力进行消费,拉取的频率需要用户自己控制(压力可控,可以一次性拉取一批数据,效率更高),拉取频繁容易造成服务端和客户端的网络传输压力,拉取间隔长又容易造成消费不及时(实时性不强
DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("test-consumer-group");
DefaultMQPushConsumer consumer = new DefaultMQPullConsumer("test-consumer-group");

Push模式也是基于pull模式的(不管是push还是pull,实际底层都是pull),只能客户端内部封装了api(每隔一段时间去pull一次)

  • 一般场景下,上游消息生产量小或者均速的时候,选择push模式
  • 在特殊场景下,例如电商大促,抢优惠券等场景可以选择pull模式

1.4 Rocket安装

二、RocketMQ 整体架构、消息模型与收发链路

RocketMQ 是一个分布式消息中间件。它通过 Producer、Consumer、NameServer、Broker 等核心组件,实现消息的发送、存储、路由、消费、重试和高可用。在 RocketMQ 5.x 的新客户端体系中,还可能通过 Proxy 作为统一接入层。

理解 RocketMQ,可以先记住一条主线:

Producer 从 NameServer 获取 Topic 路由
        ↓
选择 Topic 下的某个 MessageQueue
        ↓
将消息发送到 MessageQueue 所在的 Broker
        ↓
Broker 持久化保存消息
        ↓
Consumer 从 NameServer 获取 Topic 路由
        ↓
ConsumerGroup 内的消费者分配 MessageQueue
        ↓
Consumer 从对应 Broker 拉取并处理消息

更简洁地说:

Producer 查询路由 → 消息写入 Broker → Consumer 查询路由 → 从 Broker 拉取消息

其中:

Topic:消息的业务分类
MessageQueue:Topic 下的分区,是发送和消费负载均衡的基本单位
Broker:真正保存和投递消息的服务器
NameServer:保存 Topic 与 Broker 之间的路由关系
ConsumerGroup:一组执行相同消费逻辑的消费者
Tag:Topic 内的二级分类
Key:消息的业务检索标识

2.1 RocketMQ 的整体架构

RocketMQ 的核心组件包括:

Producer:发送消息
Consumer:消费消息
NameServer:提供路由注册与发现
Broker:存储、投递和管理消息
Proxy:RocketMQ 5.x 新客户端的接入层

整体架构如下:

                         ┌────────────────┐
                         │   NameServer   │
                         │ 保存路由信息    │
                         └───────▲────────┘
                                 │ Broker 注册
                                 │ 客户端查询路由
                 ┌───────────────┴────────────────┐
                 │                                │
        ┌────────┴────────┐              ┌────────┴────────┐
        │    Broker A     │              │    Broker B     │
        │ Topic / Queue   │              │ Topic / Queue   │
        │ 消息存储与投递    │              │ 消息存储与投递    │
        └──────▲─────┬────┘              └──────▲─────┬────┘
               │     │                           │     │
             发送   拉取消费                    发送   拉取消费
               │     │                           │     │
          Producer Consumer                 Producer Consumer

Broker 启动后,会主动向所有 NameServer 注册自己的地址、Broker 名称、Topic、MessageQueue 和读写状态等信息,并定期发送心跳。

Producer 和 Consumer 不会让 NameServer 转发每一条消息。它们只从 NameServer 获取路由信息,得到 Broker 地址后,便直接与 Broker 通信。因此 NameServer 更像一张“路由地图”,而 Broker 才是真正存储消息的“仓库”。

各组件之间的关系如下:

关系 说明
Broker → NameServer Broker 启动时注册路由,并定期发送心跳
Producer → NameServer Producer 查询 Topic 对应的 Broker 和 MessageQueue,并缓存路由
Producer → Broker Producer 根据路由选择 MessageQueue,直接向对应 Broker 发送消息
Consumer → NameServer Consumer 查询 Topic 所在的 Broker 和 MessageQueue,并缓存路由
Consumer → Broker Consumer 根据队列分配结果,直接从对应 Broker 拉取消息
Broker 包含 Topic 和 Queue 一个 Broker 可以承载多个 Topic,每个 Topic 在该 Broker 上可以配置多个 Queue
Topic 可以跨 Broker 同一个 Topic 的 MessageQueue 可以分布在多个 Broker 上,实现水平扩展
ConsumerGroup 分配 Queue 同组消费者共同分摊 MessageQueue,实现消费负载均衡

2.2 Topic、MessageQueue 与 Broker 的关系

Topic、MessageQueue 和 Broker 分别属于不同层次的概念:

Topic:业务上的消息分类
MessageQueue:Topic 下的逻辑分区
Broker:实际存储和提供消息读写服务的服务器

一个 Topic 通常包含多个 MessageQueue,这些 MessageQueue 可以分布在多个 Broker 上;同时,一个 Broker 也可以承载多个 Topic 的 MessageQueue。因此,Topic 与 Broker 之间是多对多关系,而两者之间真正的连接单位是 MessageQueue

例如,order-event-topic 部署在 Broker A 和 Broker B 上,每个 Broker 为该 Topic 提供 3 个 MessageQueue:

Topic:order-event-topic
│
├── Broker A
│   ├── MessageQueue(topic=order-event-topic, broker=A, queueId=0)
│   ├── MessageQueue(topic=order-event-topic, broker=A, queueId=1)
│   └── MessageQueue(topic=order-event-topic, broker=A, queueId=2)
│
└── Broker B
    ├── MessageQueue(topic=order-event-topic, broker=B, queueId=0)
    ├── MessageQueue(topic=order-event-topic, broker=B, queueId=1)
    └── MessageQueue(topic=order-event-topic, broker=B, queueId=2)

从逻辑上看,这个 Topic 一共有 6 个可供发送和消费的 MessageQueue。但不同 Broker 上的 QueueId 可以相同,因为一个 MessageQueue 的完整身份并不只是 QueueId,而是:

Topic + BrokerName + QueueId

例如,下面两个 MessageQueue 并不是同一个 Queue:

order-event-topic + Broker A + QueueId 0
order-event-topic + Broker B + QueueId 0

因此,不能简单理解为“一个 Topic 属于某一台 Broker”。更准确的关系是:

一个 Topic 包含多个 MessageQueue;
一个 Topic 的 MessageQueue 可以分布在多个 Broker 上;
一个 Broker 可以承载多个 Topic 的 MessageQueue。

例如:

Broker A
├── order-event-topic
│   ├── Queue 0一个 Topic 下面有多个 MessageQueue(默认是 4 个)。生产者发消息时,会根据负载均衡策略(比如轮询、哈希取模)选择其中一个 Queue 来发送。
│   ├── Queue 1
│   └── Queue 2
│
└── payment-event-topic
    ├── Queue 0
    └── Queue 1

Broker B
├── order-event-topic
│   ├── Queue 0
│   ├── Queue 1
│   └── Queue 2
│
└── file-event-topic
    ├── Queue 0
    └── Queue 1

这里体现了两层关系:

从 Topic 角度:
一个 Topic 可以跨多个 Broker 部署。

从 Broker 角度:
一个 Broker 可以存储多个 Topic 的部分 MessageQueue。

这种设计主要带来两个能力。

第一,水平扩展。 一个 Topic 的 MessageQueue 分散在多个 Broker 上后,Producer 的写请求和 Consumer 的读请求可以由多个 Broker 共同承担。增加 Broker 和 MessageQueue,可以提升整个 Topic 的吞吐能力。

Producer
   ├── 消息 1 → Broker A / Queue 0
   ├── 消息 2 → Broker B / Queue 1
   └── 消息 3 → Broker A / Queue 2

第二,并行消费。 同一 ConsumerGroup 中的多个消费者可以分别处理不同的 MessageQueue,从而并行消费消息。

Consumer 1 → Broker A / Queue 0
Consumer 2 → Broker A / Queue 1
Consumer 3 → Broker B / Queue 0
Consumer 4 → Broker B / Queue 1

因此,MessageQueue 不仅是 Topic 下的逻辑分区,也是消息发送路由、消费负载均衡和局部顺序控制的基本单位。

不过,“消息存放在 MessageQueue 中”主要是业务和逻辑层面的理解。RocketMQ 在物理存储上并不会为每个 MessageQueue 单独创建一份完整的消息文件。

消息写入 Broker 后,主要涉及以下两个结构:

CommitLog:
保存完整消息内容,Broker 收到的消息通常按顺序追加写入。

ConsumeQueue:
保存某个 Topic、某个 Queue 对应的逻辑索引,
索引中记录消息在 CommitLog 中的位置、大小等信息。

可以将它们理解为:

CommitLog:真正存放消息正文的总账本
ConsumeQueue:每个 MessageQueue 对应的目录或索引

例如,Broker A 收到来自不同 Topic 和 Queue 的消息:

CommitLog

位置 100:order-event-topic / Queue 0 / 消息 A
位置 200:payment-event-topic / Queue 1 / 消息 B
位置 300:order-event-topic / Queue 1 / 消息 C
位置 400:order-event-topic / Queue 0 / 消息 D

对应的 ConsumeQueue 只保存索引:

ConsumeQueue:order-event-topic / Queue 0
├── 指向 CommitLog 位置 100
└── 指向 CommitLog 位置 400

ConsumeQueue:order-event-topic / Queue 1
└── 指向 CommitLog 位置 300

ConsumeQueue:payment-event-topic / Queue 1
└── 指向 CommitLog 位置 200

因此可以从两个视角理解 RocketMQ 的存储关系:

业务视角:

Topic
  ↓
多个 MessageQueue
  ↓
每个 MessageQueue 中存在一系列消息
物理存储视角:

完整消息统一顺序写入 CommitLog
  ↓
每个 Topic + Queue 对应一个 ConsumeQueue
  ↓
ConsumeQueue 保存指向 CommitLog 的逻辑索引

最终可以概括为:

Topic 是逻辑上的业务分类;
MessageQueue 是 Topic 下的逻辑分区;
Broker 是 MessageQueue 的实际承载节点;
Topic 与 Broker 是通过 MessageQueue 建立的多对多关系;
消息正文实际存储在 Broker 的 CommitLog 中;
ConsumeQueue 为每个 Topic-Queue 保存消息索引。

2.3 Topic:消息的业务分类

在 RocketMQ 中,Topic 是一个逻辑概念,MessageQueue 是物理概念。 打个比方:Topic 就像是一个“主题文件夹”,而 MessageQueue 是这个文件夹下面的“物理文件”

  • Topic:生产者发送消息时指定 Topic,消费者订阅时指定 Topic。它是消息的分类标签。
  • MessageQueue:Topic 下面实际存储消息的最小物理单元。一条消息最终会被写到某个具体的 MessageQueue 里。 一个 Topic 下面有多个 MessageQueue(默认是 4 个)。生产者发消息时,会根据负载均衡策略(比如轮询、哈希取模)选择其中一个 Queue 来发送。

Topic 表示一类业务消息,例如:

order-event-topic
file-event-topic
payment-event-topic
notification-topic

Topic 是生产者和消费者之间的逻辑连接点:

Producer 向 Topic 发送消息
Consumer 订阅 Topic 消费消息

Producer 并不知道消息最终由哪个具体 Consumer 实例处理,Consumer 也不需要知道消息来自哪个业务服务实例。这种基于 Topic 的通信方式实现了生产者与消费者之间的解耦。

企业项目中通常按照稳定的业务领域设计 Topic。例如订单创建、订单支付和订单取消可以放在同一个订单事件 Topic 中,再通过 Tag 区分事件类型。

不建议:

每个接口创建一个 Topic
每个用户创建一个 Topic
把完全无关的业务消息放进同一个 Topic

2.4 MessageQueue:Topic 下的分区

这是 RocketMQ 最精妙的设计之一,也是它为什么能具备极高写入吞吐量的核心机密。RocketMQ 采用了 CommitLog + ConsumeQueue 的两层存储结构。

  • CommitLog(仓库):所有消息顺序写入同一个大文件(每个 Broker 只有一个 CommitLog)。无论消息属于哪个 Topic,都顺序追加到 CommitLog 里——这种顺序写入的方式让 RocketMQ 拥有了极高的写入吞吐量(因为磁盘顺序写远快于随机写)。
  • ConsumeQueue(货架标签):每个 MessageQueue 对应一个 ConsumeQueue 文件,它相当于一个“索引”,里面记录了每条消息在 CommitLog 中的物理偏移量、大小、Tag 哈希值等信息。

Consumer 消费时,先查 ConsumeQueue 找到消息在 CommitLog 里的位置,再去 CommitLog 读取消息内容。

一个 Topic 通常包含多个 MessageQueue:

order-event-topic
├── Queue 0
├── Queue 1
├── Queue 2
└── Queue 3

Producer 发送消息时,会选择其中一个 MessageQueue。ConsumerGroup 中的消费者实例,则会共同分配这些 MessageQueue。

MessageQueue 主要决定两个能力。

1. 提高并发吞吐量

假设一个 Topic 有 4 个 Queue,同一个 ConsumerGroup 中启动了 3 个消费者:

Consumer 1 → Queue 0、Queue 3
Consumer 2 → Queue 1
Consumer 3 → Queue 2

三个消费者可以并行处理不同 Queue 中的消息。多个 Queue 意味着多个消费者可以并行消费,吞吐量成倍提升。如果 Topic 只有一个 Queue,即使启动多个消费者,同一个 ConsumerGroup 中通常也只有一个消费者能够处理该 Queue,其余实例无法获得额外并行能力。因此,在集群消费模式下:同一个 ConsumerGroup 的有效消费并行度,通常不会超过 MessageQueue 数量。

2. 实现顺序消息

RocketMQ 的顺序消息,保证的是同一个 Queue 内的消息严格有序。所以如果你要保证某一批消息的顺序(比如同一个订单的创建→支付→发货),只需要把它们的 Key(比如订单 ID)用哈希取模分配到同一个 Queue 即可。

RocketMQ 可以保证同一个 MessageQueue 内消息的存储顺序,但不能天然保证一个 Topic 中多个 Queue 之间的全局顺序。 例如同一订单的事件:

ORDER_CREATED
ORDER_PAID
ORDER_SHIPPED

如果三条消息分别进入不同 Queue,多个消费者并行消费时,执行顺序可能发生变化。因此顺序消息通常需要根据业务键key选择 Queue:

同一个 orderId
      ↓
始终选择同一个 MessageQueue
      ↓
同一订单的事件保持局部顺序

需要同时满足:

发送端:同一业务键选择同一个 Queue
Broker:该 Queue 内顺序存储
消费端:对该 Queue 顺序消费

2.5 Tag:Topic 内的二级分类

Tag 用于在同一个 Topic 中进一步区分不同类型的消息。

例如:

Topic:order-event-topic

Tag:
ORDER_CREATED
ORDER_PAID
ORDER_CANCELLED
ORDER_FINISHED

Consumer 在订阅 Topic 时,可以用 Tag 来过滤消息,只消费自己感兴趣的那一类。 Tag 的使用方式

  • 生产者发送时指定 Tag:SendResult result = producer.send(msg, "order_create");
  • 消费者订阅时指定 Tag:consumer.subscribe("order_topic", "order_create || order_pay");

Topic 和 Tag 的关系可以理解为:

Topic:大的业务领域
Tag:该领域中的具体事件类型

Producer 发送消息时设置 Tag,Consumer 订阅时设置过滤表达式。Broker 会根据订阅条件过滤消息,再向消费者提供符合条件的消息。 Tag 适合简单、稳定的分类条件。如果过滤逻辑更复杂,可以为消息设置自定义属性,并通过 SQL92 表达式过滤。 不建议把高基数、频繁变化的业务字段作为 Tag,例如:

userId
orderId
fileId
时间戳

这些字段更适合作为消息属性或 Key。

2.6 Key:消息的业务检索标识

Key 是消息的业务检索标识,通常设置为:

订单号
文件 ID
支付流水号
事件 ID

例如:

Topic = order-event-topic
Tag   = ORDER_CREATED
Key   = order-10001

Key 主要用于消息查询和故障排查,并不会直接决定消息进入哪个 MessageQueue。

1. 消息查询:RocketMQ 提供了根据 Key 查询消息的功能。当你在控制台输入订单号查询消息时,背后利用的就是 Key 索引。这个索引的核心是Hash 索引——消息发送时,Broker 会根据 Key 的哈希值构建索引条目,存储在索引文件中,查询时通过哈希快速定位到具体的消息。

2. 幂等处理:Consumer 可以根据 Key 来做幂等判断,防止消息重复消费导致业务数据出错。

Key 的设计建议

  • 尽量用业务上唯一且不变的字段(如订单 ID、用户 ID)
  • 如果一条消息包含多个业务维度,选最核心的那个作为 Key
  • 不建议把时间戳、随机数这类变化频繁的值作为 Key(无法精准定位)

💡 小贴士:RocketMQ 的消息体最大支持 4MB(5.x 版本可通过配置调整),但如果你的消息体很大(比如超过 1MB),建议考虑压缩或把大内容存到 OSS,消息里只存引用路径。

Broker 可以为消息 Key 建立索引。当出现消息未消费、重复消费或业务状态不一致时,可以根据 Key 在 RocketMQ 控制台或管理工具中查询消息轨迹和存储信息。

企业中通常建议保存以下标识:

eventId:事件唯一标识,主要用于消费幂等
businessId:订单号、文件 ID 等业务主键
traceId:整个调用链的追踪标识
Key:RocketMQ 中用于查询消息的索引标识

这些字段可以部分重合。例如,可以直接把 eventIdbusinessId 设置为 Key。

2.7 ConsumerGroup:一类消费行为

ConsumerGroup 不是一台消费者,而是一组具有相同消费逻辑的消费者实例。 例如:

Topic:order-event-topic

ConsumerGroup:inventory-service-group
├── inventory-service-1
├── inventory-service-2
└── inventory-service-3

ConsumerGroup:notification-service-group
├── notification-service-1
└── notification-service-2

库存服务的三个实例属于同一个 ConsumerGroup,它们共同分摊该 Topic 的 MessageQueue。

通知服务使用另一个 ConsumerGroup,因此同一条订单消息既可以被库存服务消费,也可以被通知服务消费。

需要记住两条规则:

同一个 ConsumerGroup 内:
多个实例负载均衡,一条消息通常只由其中一个实例处理

不同 ConsumerGroup 之间:
消费进度相互独立,每个消费组都可以消费一遍消息

例如一条订单创建消息:

inventory-service-group
    由组内某一个库存服务实例消费

notification-service-group
    由组内某一个通知服务实例消费

search-service-group
    由组内某一个搜索服务实例消费

因此,同一个 Topic 可以被多个不同业务系统订阅,但不同业务必须使用不同的 ConsumerGroup。

如果库存服务和通知服务错误地使用相同 ConsumerGroup,它们会相互竞争消息,导致库存服务只收到一部分消息,通知服务也只收到一部分消息。

同一 ConsumerGroup 中的消费者通常还应保持一致的:

订阅 Topic
Tag 或过滤表达式
消费模式
消息处理语义

否则可能出现消费行为不一致或部分消息无法按预期处理。

2.8 Producer:消息生产者

Producer 通常集成在业务服务中,负责将业务事件封装成消息并发送到 Topic。 基本使用步骤是:

1. 创建 Producer,并指定 ProducerGroup
2. 配置 NameServer 或接入地址
3. 启动 Producer
4. 创建消息,设置 Topic、Tag、Key、属性和消息体
5. 发送消息
6. 应用关闭时关闭 Producer

例如订单创建成功后发送消息:

{
  "eventId": "evt-20260711-0001",
  "orderId": 10001,
  "userId": 20001,
  "eventType": "ORDER_CREATED",
  "occurredAt": "2026-07-11T20:30:00"
}

Producer 发送消息时,并不是直接指定 Consumer,而是指定:

Topic
Tag
Key
消息属性
消息体

典型发送过程如下:

Producer 查询 Topic 路由
        ↓
获得 Broker 和 MessageQueue 列表
        ↓
根据负载均衡策略选择一个 Queue
        ↓
将消息发送到该 Queue 所在 Broker
        ↓
Broker 写入消息
        ↓
Broker 返回发送结果

RocketMQ 客户端会缓存 Topic 路由,并定期刷新,因此不是每发送一条消息都访问一次 NameServer。

RocketMQ 常见发送方式包括:

发送方式 特点 适用场景
同步发送 等待 Broker 返回结果,调用线程阻塞 订单、支付等需要明确发送结果的消息
异步发送 立即返回,通过回调接收发送结果 高吞吐且不能长时间阻塞业务线程
单向发送 只发送,不等待 Broker 返回结果 日志、监控等允许少量丢失的场景
顺序发送 根据业务键选择固定 Queue 同一订单、同一账户的有序事件
批量发送 一次发送多条消息 批量日志、批量数据同步
事务消息 本地事务与消息发送保持最终一致 订单创建、支付结果、库存操作等

2.9 Consumer:消息消费者

Consumer 订阅 Topic,从 Broker 获取消息并执行实际业务,例如:

扣减库存
创建物流任务
更新搜索索引
发送短信
清理文件
生成缩略图

基本使用步骤是:

1. 创建 Consumer,并指定 ConsumerGroup
2. 配置 NameServer 或接入地址
3. 订阅 Topic 和 Tag
4. 注册消息监听器
5. 启动 Consumer
6. 接收并处理消息

RocketMQ 消费者底层采用 Pull 模式,即客户端主动向 Broker 拉取消息。传统 PushConsumer 所谓的 Push,并不是 Broker 主动向 Consumer 建立推送,而是客户端在 Pull 基础上进行了长轮询、缓存和监听器回调封装,使开发者感受起来像“消息主动推送”。

可靠消费过程如下:

Consumer 从 Broker 拉取消息
        ↓
执行本地业务逻辑
        ↓
业务处理成功
        ↓
返回消费成功
        ↓
更新消费进度

如果消费代码抛出异常或返回失败:

消费失败
    ↓
Broker 或客户端按照重试策略重新投递
    ↓
多次失败后进入死信队列

因此,“Consumer 收到消息”并不代表业务已经处理完成。只有业务执行成功,并正确返回成功状态,才能认为本次消费完成。

业务代码还必须考虑重复消费,因为在网络异常、消费超时、ACK 丢失、Consumer 重启等情况下,同一条消息可能被多次投递。

常见幂等设计包括:

基于 eventId 建唯一索引
基于业务状态进行条件更新
使用去重表记录已消费事件
使用 Redis SETNX 做短期去重
利用数据库唯一约束阻止重复写入

2.10 NameServer:路由注册与发现中心

NameServer 是 RocketMQ 的“注册中心”和“路由字典”。 它不存消息(消息存 CommitLog),它的作用是:告诉生产者和消费者,Topic 的 Queue 到底在哪台 Broker 上

NameServer 是 RocketMQ 的轻量级路由注册与发现中心,主要保存:

有哪些 Broker
Broker 的地址和角色
每个 Topic 位于哪些 Broker
每个 Topic 有哪些 MessageQueue
Broker 的读写状态
Broker 的存活情况

Broker 启动后,会向所有 NameServer 注册路由并定时发送心跳。

心跳是一种“我还活着”的定时保活机制。 就像你上班打卡一样,Broker 每隔一段时间(默认 30 秒)主动给 NameServer 发一个 “心跳包”(一个极小的网络数据包),告诉它:

“NameServer 大哥,我还活着,我的路由信息没变(或变更如下)。”

心跳的核心作用:

作用 解释
服务发现 NameServer 通过心跳知道当前集群有哪些 Broker 在线上。
故障剔除 如果 NameServer 超过 120 秒(默认) 没收到某台 Broker 的心跳,就认为它挂了,立即从路由表中移除该 Broker 的所有 Queue。
路由更新 如果 Broker 上的 Topic 新增了 Queue,也会通过心跳同步给 NameServer。

Producer 和 Consumer 从 NameServer 查询 Topic 路由,随后直接与 Broker 通信:

Producer / Consumer
        ↓ 查询路由
NameServer
        ↓ 返回 Broker 地址和 Queue 信息
Producer / Consumer
        ↓ 直接通信
Broker

NameServer 不负责:

保存业务消息
转发每一条消息
执行消息消费
保存完整业务消费状态

客户端会缓存路由信息。因此,NameServer 短时间不可用时,已经获取过路由的 Producer 和 Consumer 可能仍能继续与原 Broker 通信。

但是以下能力会受到影响:

首次查询新 Topic
发现新加入的 Broker
感知 Broker 路由变化
定期刷新路由

生产环境通常部署多个 NameServer。各 NameServer 节点相互独立,Broker 会向每一个 NameServer 注册,而不是依靠 NameServer 之间的数据同步。

2.11 Broker:消息存储与投递核心

Broker 是 RocketMQ 最核心的服务端组件,主要负责:

接收 Producer 消息
将消息写入 CommitLog
维护 ConsumeQueue 和 IndexFile
管理 Topic 和 MessageQueue
响应 Consumer 拉取请求
保存和管理消费进度
执行消息重试和死信处理
处理延时消息
处理事务消息回查
向 NameServer 注册路由
执行主从复制
参与故障恢复和主从切换

Broker 才是消息真正存储的位置。

可以使用一个简单类比:

NameServer:地图或导航系统
Broker:消息仓库
Topic:仓库中的业务分类
MessageQueue:分类下的逻辑分区
Producer:入库方
Consumer:出库并处理消息的一方

Broker Master 通常负责消息读写,Slave 从 Master 同步数据。在传统主从架构下,Slave 不能接收 Producer 写入,但可以在特定条件下承担读请求。

2.12 Proxy:RocketMQ 5.x 的统一接入层

RocketMQ 5.x 引入了新的 gRPC 客户端体系,客户端通常通过 Proxy 访问 RocketMQ 服务。

Proxy 可以采用两种部署方式:

Local 模式:
Proxy 与 Broker 同节点或同进程部署

Cluster 模式:
Proxy 独立组成集群,作为统一接入层

Proxy 偏无状态,主要负责:

接收 gRPC 客户端请求
协议转换
客户端连接管理
路由和请求转发
为多语言客户端提供统一接入

企业项目中需要先确认使用的是哪套客户端。

传统 RocketMQ Java Client
    使用 Remoting 协议
    Producer 和 Consumer 通常直接访问 Broker
    Spring 项目常见 rocketmq-spring-boot-starter

RocketMQ 5.x gRPC Client
    客户端通常访问 Proxy
    使用 Endpoint 或 Proxy 地址
    更适合统一的多语言客户端体系

需要注意:

RocketMQ 服务端版本是 5.x
不代表项目必须使用 5.x gRPC Client

很多 RocketMQ 5.x 集群仍然兼容传统 Java Client。传统客户端和新 gRPC 客户端在 API、接入地址、协议和部分行为上存在差异,配置时不能混用。

2.13 一条消息的完整发送与消费链路

以订单创建事件为例,完整过程如下。

第一步:Broker 注册路由

Broker 启动后向所有 NameServer 注册:

Broker 地址
Broker 名称
Master / Slave 信息
Topic 列表
MessageQueue 数量
读写权限

第二步:Producer 获取 Topic 路由

Producer 准备向 order-event-topic 发送消息时,从 NameServer 查询:

该 Topic 位于哪些 Broker
每个 Broker 有哪些可写 MessageQueue

客户端将路由缓存到本地,并定期刷新。

第三步:Producer 选择 MessageQueue

普通消息通常通过负载均衡选择 Queue:

Queue 0
Queue 1
Queue 2
Queue 3

顺序消息则根据 orderId 等业务键固定选择 Queue。

第四步:消息写入 Broker

Producer 直接连接 MessageQueue 所在的 Broker,Broker 接收消息并写入存储系统。

从逻辑上看:

消息进入某个 MessageQueue

从物理存储上看:

消息主体顺序写入 CommitLog
对应 ConsumeQueue 写入逻辑索引

第五步:Consumer 获取路由和 Queue

Consumer 从 NameServer 获取 Topic 路由,并根据 ConsumerGroup 内的实例情况执行 Queue 分配。

例如:

Topic 有 4 个 Queue
ConsumerGroup 有 2 个实例

Consumer 1 → Queue 0、Queue 2
Consumer 2 → Queue 1、Queue 3

第六步:Consumer 拉取并处理消息

Consumer 根据分配结果,向对应 Broker 拉取消息并执行本地业务。

Broker 返回消息
        ↓
Consumer 执行业务
        ↓
成功:确认消费并推进消费进度
失败:触发后续重试

第七步:主从复制

Master 接收消息后,根据配置将消息复制给 Slave,以提高数据可靠性和读取可用性。

完整链路可以概括为:

Broker 向 NameServer 注册 Topic 路由
                ↓
Producer 查询 Topic 路由
                ↓
Producer 选择 MessageQueue
                ↓
Producer 直接向 Broker 发送消息
                ↓
Broker 持久化消息并进行主从复制
                ↓
Consumer 查询 Topic 路由
                ↓
ConsumerGroup 分配 MessageQueue
                ↓
Consumer 从 Broker 拉取消息
                ↓
业务成功后确认消费
                ↓
失败则进入重试与死信机制

2.14 核心概念关系总结

Producer
    向 Topic 发送消息
        ↓
Topic
    表示业务消息类别
        ↓
MessageQueue
    是 Topic 下的分区,决定并行度和局部顺序
        ↓
Broker
    保存 MessageQueue 对应的消息数据
        ↓
ConsumerGroup
    订阅 Topic,并在组内分配 MessageQueue
        ↓
Consumer
    拉取消息并执行实际业务

Tag 和 Key 则描述一条消息的附加语义:

Tag:
Topic 内的二级分类,用于订阅过滤

Key:
业务检索标识,用于查询和问题排查

最终可以记住:

Topic 决定消息属于哪类业务;
MessageQueue 决定消息写到哪个分区以及如何并行消费;
Broker 决定消息实际存在哪里;
Tag 决定消费者关注哪些事件类型;
Key 决定出现问题时如何检索消息;
ConsumerGroup 决定消息由哪一类业务消费;
Consumer 实例决定最终由哪台机器执行消息处理。

三、RocketMQ设计与原理

1 消息存储

Broker内部结构

消息存储是RocketMQ中最为复杂和最为重要的一部分,本节将分别从RocketMQ的消息存储整体架构、PageCache与Mmap内存映射以及RocketMQ中两种不同的刷盘方式三方面来分别展开叙述。

1.1 消息存储整体架构

RocketMQ 经典本地文件存储主要涉及 CommitLogConsumeQueueIndexFile 三类结构。三者并不是分别保存三份完整消息,而是采用“消息主体与索引分离”的设计:完整消息统一写入 CommitLog,ConsumeQueue 和 IndexFile 保存指向 CommitLog 的索引。消息存储架构图中主要有下面三个跟消息存储相关的文件构成。

(1) CommitLog:消息主体以及元数据的存储主体,存储 Producer 端写入的消息主体内容,单条消息记录不是定长的。单个文件大小默认 1G, 文件名长度为20位,左边补零,剩余为起始偏移量,比如00000000000000000000代表了第一个文件,起始偏移量为0,文件大小为1G=1073741824;当第一个文件写满了,第二个文件为00000000001073741824,起始偏移量为1073741824,以此类推。消息主要采用顺序写入日志文件的方式,当当前文件写满后,再写入下一个文件

(2) ConsumeQueue:消息消费队列,引入的目的主要是提高消息消费的性能,由于RocketMQ是基于主题topic的订阅模式,消息消费是针对主题进行的,如果直接遍历 CommitLog 并根据 Topic 检索消息,效率会非常低。Consumer 可以根据 ConsumeQueue 快速定位待消费消息。其中,ConsumeQueue(逻辑消费队列)作为消费消息的索引,保存了指定Topic下的队列消息在CommitLog中的起始物理偏移量offset,消息大小size和消息Tag的HashCode值。ConsumeQueue 文件可以看成基于 Topic 和 Queue 的 CommitLog 逻辑索引文件,故consumequeue文件夹的组织方式如下:topic/queue/file三层组织结构,具体存储路径为:$HOME/store/consumequeue/{topic}/{queueId}/{fileName}。同样consumequeue文件采取定长设计,每一个条目共20个字节,分别为8字节的commitlog物理偏移量、4字节的消息长度、8字节tag hashcode,单个文件由30W个条目组成,可以像数组一样随机访问每一个条目,每个ConsumeQueue文件大小约5.72M;

(3) IndexFile:IndexFile(索引文件)提供了一种可以通过key或时间区间来查询消息的方法。Index文件的存储位置是:$HOME/store/index/{fileName},文件名fileName是以创建时的时间戳命名的,经典 4.x 实现中,单个 IndexFile 文件大小约为 400M,一个 IndexFile 可以保存约 2000W 个索引,IndexFile 的底层结构类似在文件中实现的 HashMap,故RocketMQ的索引文件其底层实现为hash索引。

在上面的RocketMQ的消息存储整体架构图中可以看出,RocketMQ采用的是混合型的存储结构,即为Broker单个实例下所有的队列共用一个日志数据文件(即为CommitLog)来存储。RocketMQ的混合型存储结构(多个Topic的消息实体内容都存储于一个CommitLog中)针对Producer和Consumer分别采用了数据和索引部分相分离的存储结构,Producer发送消息至Broker端,然后Broker端使用同步或者异步的方式对消息刷盘持久化,保存至CommitLog中。消息刷盘持久化到 CommitLog 后,可以避免 Broker 进程崩溃或机器重启导致的数据丢失;但磁盘损坏等极端故障仍需要主从复制或副本机制保障。只要消息仍然存在,Consumer 就有机会继续消费。当无法拉取到消息后,可以等下一次消息拉取,同时服务端也支持长轮询模式,如果一个消息拉取请求未拉取到消息,Broker 可以在配置的长轮询等待时间内挂起请求(经典实现常见为约 15s),只要这段时间内有新消息到达,将直接返回给消费端。这里,RocketMQ的具体做法是,使用Broker端的后台服务线程—ReputMessageService持续转发 CommitLog 中的新消息,并异步构建 ConsumeQueue(逻辑消费队列)和 IndexFile(索引文件)

1.2 页缓存与内存映射

页缓存(PageCache)是OS对文件的缓存,用于加速对文件的读写。一般来说,顺序文件读写在命中 PageCache 时可以获得非常高的性能,但不能简单等同于内存读写速度,主要原因就是由于OS使用PageCache机制对读写访问操作进行了性能优化,将一部分的内存用作PageCache。对于数据的写入,OS会先写入至Cache内,随后由操作系统的后台回写机制将脏页异步刷入物理磁盘。对于数据的读取,如果一次读取文件时出现未命中PageCache的情况,OS从物理磁盘上访问读取文件的同时,会顺序对其他相邻块的数据文件进行预读取。

在RocketMQ中,ConsumeQueue逻辑消费队列存储的数据较少,并且是顺序读取,在page cache机制的预读取作用下,Consume Queue文件的读性能几乎接近读内存,即使存在一定消息堆积,通常也能保持较好的读取性能。而对于CommitLog消息存储的日志数据文件来说,读取消息内容时候会产生较多的随机访问读取,可能增加随机 I/O 开销并影响读取性能。如果选择合适的系统IO调度算法,例如根据磁盘类型和操作系统版本选择合适的 I/O 调度策略,随机读的性能也会有所提升。

另外,RocketMQ主要通过MappedByteBuffer对文件进行读写操作。其中,利用了NIO中的FileChannel模型将磁盘上的物理文件直接映射到用户态的内存地址中(这种Mmap的方式减少了传统IO将磁盘文件数据在操作系统内核地址空间的缓冲区和用户应用程序地址空间的缓冲区之间来回进行拷贝的性能开销),将对文件的操作转化为直接对内存地址进行操作,从而极大地提高了文件的读写效率(RocketMQ 的存储文件通常采用固定大小的文件段,便于内存映射;但 CommitLog 中的单条消息记录仍然是变长的)。

1.3 消息刷盘

(1) 同步刷盘:如上图所示,只有在消息按照同步刷盘策略写入磁盘后,Broker 才会向 Producer 返回成功 ACK。同步刷盘对MQ消息可靠性来说是一种不错的保障,但是性能上会有较大影响,一般适用于对消息持久化可靠性要求较高的业务。

(2) 异步刷盘:能够充分利用OS的PageCache的优势,只要消息写入PageCache即可将成功的ACK返回给Producer端。消息刷盘采用后台异步线程提交的方式进行,降低了读写延迟,提高了MQ的性能和吞吐量。

2 通信机制

RocketMQ消息队列集群主要包括NameServer、Broker(Master/Slave)、Producer、Consumer 4 个角色,基本通讯流程如下:

(1) Broker启动后需要完成一次将自己注册至NameServer的操作;随后按照配置的周期定时向 NameServer 上报 Topic 路由和 Broker 状态信息(经典实现中通常为约 30s)。

(2) 消息生产者Producer作为客户端发送消息时候,需要根据消息的Topic从本地缓存的TopicPublishInfoTable获取路由信息。如果本地没有对应路由或路由需要刷新,就会从 NameServer 重新拉取;客户端也会按照配置周期定时更新路由信息(经典实现中通常为约 30s)。

(3) 消息生产者Producer根据2)中获取的路由信息选择一个队列(MessageQueue)进行消息发送;Broker作为消息的接收者接收消息并落盘存储。

(4) 消息消费者Consumer根据2)中获取的路由信息,并在完成客户端负载均衡后,选择其中的某一个或者某几个消息队列来拉取消息并进行消费。

从上面 1)~4)可以看出,Producer、Consumer、Broker 和 NameServer 之间都会发生通信(这里只说了MQ的部分通信),因此如何设计一个良好的网络通信模块在MQ中至关重要,它将决定RocketMQ集群整体的消息传输能力与最终的性能。

rocketmq-remoting 模块是 RocketMQ消息队列中负责网络通信的模块,它几乎被其他所有需要网络通信的模块(诸如rocketmq-client、rocketmq-broker、rocketmq-namesrv)所依赖和引用。为了实现客户端与服务器之间高效的数据请求与接收,RocketMQ 自定义了通信协议,并在 Netty 基础上实现了 Remoting 通信模块。

2.1 Remoting通信类结构

2.2 协议设计与编解码

在Client和Server之间完成一次消息发送时,需要对发送的消息进行一个协议约定,因此就有必要自定义RocketMQ的消息协议。同时,为了高效地在网络中传输消息和对收到的消息读取,就需要对消息进行编解码。在RocketMQ中,RemotingCommand 是消息传输过程中请求和响应的统一封装,包含请求码、扩展字段、消息体以及编解码相关信息。

Header字段 类型 Request说明 Response说明
code int 请求操作码,应答方根据不同的请求码进行不同的业务处理 应答响应码。0表示成功,非0则表示各种错误
language LanguageCode 请求方实现的语言 应答方实现的语言
version int 请求方程序的版本 应答方程序的版本
opaque int 相当于requestId,在同一个连接上的不同请求标识码,与响应消息中的相对应 应答不做修改直接返回
flag int 区分是普通RPC还是onewayRPC的标志 区分是普通RPC还是onewayRPC的标志
remark String 传输自定义文本信息 传输自定义文本信息
extFields HashMap<String, String> 请求自定义扩展信息 响应自定义扩展信息

可见传输内容主要可以分为以下4部分:

(1) 消息长度:总长度,四个字节存储,占用一个int类型;

(2) 序列化类型&消息头长度:同样占用一个int类型,第一个字节表示序列化类型,后面三个字节表示消息头长度;

(3) 消息头数据:经过序列化后的消息头数据; (4) 消息主体数据:消息主体的二进制字节数据内容;

2.3 消息的通信方式和流程

在RocketMQ消息队列中支持通信的方式主要有同步(sync)、异步(async)、单向(oneway) 三种。其中“单向”通信模式相对简单,可用于对响应结果不敏感的场景,例如部分心跳或日志类请求,无需关注其Response。这里,主要介绍RocketMQ的异步通信流程。

2.4 Reactor多线程设计

RocketMQ的RPC通信采用Netty组件作为底层通信库,同样也遵循了Reactor多线程模型,同时又在这之上做了一些扩展和优化。

上面的框图中可以大致了解RocketMQ中NettyRemotingServer的Reactor 多线程模型。一个 Reactor 主线程(eventLoopGroupBoss,即为上面的1)负责监听 TCP网络连接请求,建立好连接,创建SocketChannel,并注册到selector上。RocketMQ 可以根据操作系统和配置选择 NIO 或 Epoll,然后监听和处理网络数据。拿到网络数据后,再丢给Worker线程池(eventLoopGroupSelector,即为上面的“N”,源码中默认设置为3),在真正执行业务逻辑之前需要进行SSL验证、编解码、空闲检查、网络连接管理,这些工作交给defaultEventExecutorGroup(即为上面的“M1”,源码中默认设置为8)去做。而处理业务操作放在业务线程池中执行,根据 RemotingCommand 的业务请求码 code去processorTable这个本地缓存变量中找到对应的 processor,然后封装成task任务后,提交给对应的业务processor处理线程池来执行(sendMessageExecutor,以发送消息为例,即为上面的 “M2”)。通过网络接入、编解码和业务处理等不同线程池进行职责隔离,可以避免耗时业务阻塞网络 I/O 线程。具体线程数量和名称会随版本及配置变化。

线程数 线程名 线程具体说明
1 NettyBoss_%d Reactor 主线程
N NettyServerEPOLLSelector_%d_%d Reactor 线程池
M1 NettyServerCodecThread_%d Worker线程池
M2 RemotingExecutorThread_%d 业务processor处理线程池

3 消息过滤

RocketMQ分布式消息队列的消息过滤方式有别于其它MQ中间件,由 Consumer 声明订阅条件,主要在 Broker 端执行过滤。RocketMQ这么做是在于其Producer端写入消息和Consumer端订阅消息采用分离存储的机制来实现的,Consumer 拉取消息时,Broker 会先通过 ConsumeQueue 获取逻辑索引并执行初步过滤,再根据索引从 CommitLog 读取真正的消息内容,因此过滤机制与存储结构密切相关。其ConsumeQueue的存储结构如下,可以看到其中有8个字节存储的Message Tag的哈希值,基于Tag的消息过滤正是基于这个字段值的。

主要支持如下2种的过滤方式

(1) Tag过滤方式:Consumer端在订阅消息时除了指定Topic还可以指定TAG,如果 Consumer 需要订阅多个 Tag,可以在订阅表达式中使用 || 分隔;一条消息通常只设置一个 Tag。其中,Consumer端会将这个订阅请求构建成一个 SubscriptionData,发送一个Pull消息的请求给Broker端。Broker端从RocketMQ的文件存储层—Store读取数据之前,会用这些数据先构建一个MessageFilter,然后传给Store。Store从 ConsumeQueue读取到一条记录后,会用它记录的消息tag hash值去做过滤,由于在服务端只是根据hashcode进行判断,无法精确对tag原始字符串进行过滤,为了避免 Hash 冲突造成误判,Broker 在读取 CommitLog 中的完整消息后还会结合原始 Tag 再次进行精确校验。

(2) SQL92的过滤方式:这种方式的大致做法和上面的Tag过滤方式一样,只是在Store层的具体过滤过程不太一样,真正的 SQL expression 的构建和执行由rocketmq-filter模块负责的。为了降低逐条执行 SQL 表达式的开销,RocketMQ 可以借助 BloomFilter 等机制进行预过滤,再对候选消息执行精确判断。SQL92的表达式上下文为消息的属性。

4 负载均衡

在传统 RocketMQ 4.x Java Client 中,Producer 发送队列选择和 Consumer Rebalance 主要在客户端完成,具体来说,主要可以分为Producer端发送消息时候的负载均衡和Consumer端订阅消息的负载均衡。

4.1 Producer的负载均衡

Producer端在发送消息的时候,会先根据Topic找到指定的TopicPublishInfo,在获取了TopicPublishInfo路由信息后,RocketMQ的客户端在默认方式下selectOneMessageQueue()方法会从TopicPublishInfo中的messageQueueList中选择一个队列(MessageQueue)进行发送消息。具体的容错策略均在MQFaultStrategy这个类中定义。这里有一个sendLatencyFaultEnable开关变量,如果开启,在随机递增取模的基础上,再过滤掉not available的Broker代理。所谓的"latencyFaultTolerance",是指对之前失败的,按一定的时间做退避。例如,如果上次请求的latency超过550L ms,就退避30000L ms;超过1000L,就退避60000L;如果关闭,采用随机递增取模的方式选择一个队列(MessageQueue)来发送消息,延迟故障规避机制是提升发送可用性的重要手段之一。

4.2 Consumer的负载均衡

在传统 RocketMQ 4.x 客户端中,Push/Pull 消费都基于拉取机制获取消息,而在Push模式只是对pull模式的一种封装,其本质实现为消息拉取线程在从服务器拉取到一批消息后,然后提交到消息消费线程池后,随后继续向服务器发起下一次拉取请求。如果未拉取到消息,则延迟一下又继续拉取。在两种基于拉模式的消费方式(Push/Pull)中,均需要Consumer端知道从Broker端的哪一个消息队列中去获取消息。因此,有必要在Consumer端来做负载均衡,即Broker端中多个MessageQueue分配给同一个ConsumerGroup中的哪些Consumer消费。

1、Consumer端的心跳包发送

在Consumer启动后,它就会通过定时任务不断地向RocketMQ集群中的所有Broker实例发送心跳包(其中包含了,消息消费分组名称、订阅关系集合、消息通信模式和客户端id的值等信息)。Broker端在收到Consumer的心跳消息后,会将它维护在ConsumerManager的本地缓存变量—consumerTable,同时并将封装后的客户端网络通道信息保存在本地缓存变量—channelInfoTable中,为后续 Consumer 端 Rebalance 提供消费者列表和连接信息。

2、Consumer端实现负载均衡的核心类—RebalanceImpl

在Consumer实例的启动流程中的启动MQClientInstance实例部分,会完成负载均衡服务线程—RebalanceService 的启动,并按照配置周期触发重平衡(经典实现中通常约为 20s)。通过查看源码可以发现,RebalanceService线程的run()方法最终调用的是RebalanceImpl类的rebalanceByTopic()方法,该方法是实现Consumer端负载均衡的核心。这里,rebalanceByTopic()方法会根据消费者通信类型为“广播模式”还是“集群模式”做不同的逻辑处理。这里主要来看下集群模式下的主要处理流程:

(1) 从rebalanceImpl实例的本地缓存变量—topicSubscribeInfoTable中,获取该Topic主题下的消息消费队列集合(mqSet);

(2) 根据topic和consumerGroup为参数调用mQClientFactory.findConsumerIdList()方法向Broker端发送获取该消费组下消费者Id列表的RPC通信请求(Broker端基于前面Consumer端上报的心跳包数据而构建的consumerTable做出响应返回,业务请求码:GET_CONSUMER_LIST_BY_GROUP);

(3) 先对Topic下的消息消费队列、消费者Id排序,然后用消息队列分配策略算法(默认为:消息队列的平均分配算法),计算出待拉取的消息队列。这里的平均分配算法,类似于分页的算法,将所有MessageQueue排好序类似于记录,将所有消费端Consumer排好序类似页数,并求出每一页需要包含的平均size和每个页面记录的范围range,最后遍历整个range而计算出当前Consumer端应该分配到的记录(这里即为:MessageQueue)。

!

(4) 然后,调用updateProcessQueueTableInRebalance()方法,具体的做法是,先将分配到的消息队列集合(mqSet)与processQueueTable做一个过滤比对。

!

  • 上图中processQueueTable标注的红色部分,表示与分配到的消息队列集合mqSet互不包含。将这些队列设置Dropped属性为true,然后查看这些队列是否可以移除出processQueueTable缓存变量,这里具体执行removeUnnecessaryMessageQueue()方法,即每隔1s 查看是否可以获取当前消费处理队列的锁,拿到的话返回true。如果等待1s后,仍然拿不到当前消费处理队列的锁则返回false。如果返回true,则从processQueueTable缓存变量中移除对应的Entry;

  • 上图中 processQueueTable 的绿色部分,表示仍属于当前分配结果的 MessageQueue,通常继续保留。如果 Push 模式下某个 ProcessQueue 长时间未拉取而过期,则会触发异常处理或重新分配,但不能把正常交集队列统一设置为 Dropped 并移除;

最后,为过滤后的消息队列集合(mqSet)中的每个MessageQueue创建一个ProcessQueue对象并存入RebalanceImpl的processQueueTable队列中(其中调用RebalanceImpl实例的computePullFromWhere(MessageQueue mq)方法获取该MessageQueue对象的下一个进度消费值offset,随后填充至接下来要创建的pullRequest对象属性中),并创建拉取请求对象—pullRequest添加到拉取列表—pullRequestList中,最后执行dispatchPullRequest()方法,将Pull消息的请求对象PullRequest依次放入PullMessageService服务线程的阻塞队列pullRequestQueue中,待该服务线程取出后向Broker端发起Pull消息的请求。

消息消费队列在同一消费组不同消费者之间的负载均衡,其核心设计理念是在一个消息消费队列在同一时间只允许被同一消费组内的一个消费者消费,一个消息消费者能同时消费多个消息队列。

5 事务消息

Apache RocketMQ在4.3.0版中已经支持分布式事务消息,这里RocketMQ 借鉴了两阶段提交的思想实现事务消息,同时增加一个补偿逻辑来处理二阶段超时或者失败的消息,如下图所示。

!

5.1 RocketMQ事务消息流程概要

上图说明了事务消息的大致方案,其中分为两个流程:正常事务消息的发送及提交、事务消息的补偿流程。

1.事务消息发送及提交:

(1) 发送消息(half消息)。

(2) 服务端响应消息写入结果。

(3) 根据发送结果执行本地事务(如果写入失败,此时half消息对业务不可见,本地逻辑不执行)。

(4) 根据本地事务状态执行Commit或者Rollback(Commit操作生成消息索引,消息对消费者可见)

2.补偿流程:

(1) 对没有Commit/Rollback的事务消息(pending状态的消息),从服务端发起一次“回查”

(2) Producer收到回查消息,检查回查消息对应的本地事务的状态

(3) 根据本地事务状态,重新Commit或者Rollback

其中,补偿阶段用于解决消息Commit或者Rollback发生超时或者失败的情况。

5.2 RocketMQ事务消息设计

1.事务消息在一阶段对用户不可见

在RocketMQ事务消息的主要流程中,一阶段的消息如何对用户不可见。其中,事务消息相对普通消息最大的特点就是一阶段发送的消息对用户是不可见的。那么,如何做到写入消息但是对用户不可见呢?RocketMQ事务消息的做法是:如果消息是half消息,将备份原消息的主题与消息消费队列,然后改变主题为RMQ_SYS_TRANS_HALF_TOPIC。由于消费组未订阅该主题,故消费端无法消费half类型的消息,然后RocketMQ会开启一个定时任务,从Topic为RMQ_SYS_TRANS_HALF_TOPIC中拉取消息进行消费,根据生产者组获取一个服务提供者发送回查事务状态请求,根据事务状态来决定是提交或回滚消息。

在RocketMQ中,消息在服务端的存储结构如下,每条消息都会有对应的索引信息,Consumer通过ConsumeQueue这个二级索引来读取消息实体内容,其流程如下:

!

RocketMQ的具体实现策略是:写入的如果事务消息,对消息的Topic和Queue等属性进行替换,同时将原来的Topic和Queue信息存储到消息的属性中,正因为消息主题被替换,故消息并不会转发到该原主题的消息消费队列,消费者无法感知消息的存在,不会消费。其实改变消息主题是RocketMQ的常用“套路”,回想一下延时消息的实现机制。

2.Commit和Rollback操作以及Op消息的引入

在完成一阶段写入一条对用户不可见的消息后,二阶段如果是Commit操作,则需要让消息对用户可见;如果是Rollback则需要撤销一阶段的消息。先说Rollback的情况。对于Rollback,本身一阶段的消息对用户是不可见的,其实不需要真正撤销消息(实际上RocketMQ也无法去真正的删除一条消息,因为是顺序写文件的)。但是区别于这条消息没有确定状态(Pending状态,事务悬而未决),需要一个操作来标识这条消息的最终状态。RocketMQ事务消息方案中引入了Op消息的概念,用Op消息标识事务消息已经确定的状态(Commit或者Rollback)。如果一条事务消息没有对应的Op消息,说明这个事务的状态还无法确定(可能是二阶段失败了)。引入Op消息后,事务消息无论是Commit或者Rollback都会记录一个Op操作。Commit相对于Rollback只是在写入Op消息前创建Half消息的索引。

3.Op消息的存储和对应关系

RocketMQ将Op消息写入到全局一个特定的Topic中通过源码中的方法—TransactionalMessageUtil.buildOpTopic();这个Topic是一个内部的Topic(像Half消息的Topic一样),不会被用户消费。Op消息的内容为对应的Half消息的存储的Offset,这样通过Op消息能索引到Half消息进行后续的回查操作。

!

4.Half消息的索引构建

在执行二阶段Commit操作时,需要构建出Half消息的索引。一阶段的Half消息由于是写到一个特殊的Topic,所以二阶段构建索引时需要读取出Half消息,并将Topic和Queue替换成真正的目标的Topic和Queue,之后通过一次普通消息的写入操作来生成一条对用户可见的消息。所以RocketMQ 事务消息在二阶段 Commit 时,会读取 Half 消息,恢复原 Topic 和 Queue 等属性,再按照普通消息写入流程生成一条对消费者可见的消息。

5.如何处理二阶段失败的消息?

如果在RocketMQ事务消息的二阶段过程中失败了,例如在做Commit操作时,出现网络问题导致Commit失败,那么需要通过一定的策略使这条消息最终被Commit。RocketMQ采用了一种补偿机制,称为“回查”。Broker端对未确定状态的消息发起回查,将消息发送到对应的Producer端(同一个Group的Producer),由Producer根据消息来检查本地事务的状态,进而执行Commit或者Rollback。Broker端通过对比Half消息和Op消息进行事务消息的回查并且推进CheckPoint(记录那些事务消息的状态是确定的)。

值得注意的是,RocketMQ 不会无限进行事务状态回查,经典版本默认最多回查 15 次;超过最大次数后,Broker 会按照实现和配置进行丢弃或回滚处理,并记录相关日志,具体行为应以所用版本为准。

6 消息查询

RocketMQ支持按照下面两种维度(“按照Message Id查询消息”、“按照Message Key查询消息”)进行消息查询。

6.1 按照MessageId查询消息

在经典 IPv4 场景下,RocketMQ 的 MessageId 通常由 16 字节信息编码而成,其中包含了消息存储主机地址(IP地址和端口),消息Commit Log offset。“按照MessageId查询消息”在RocketMQ中具体做法是:Client端从MessageId中解析出Broker的地址(IP地址和端口)和Commit Log的偏移地址后封装成一个RPC请求后通过Remoting通信层发送(业务请求码:VIEW_MESSAGE_BY_ID)。Broker端走的是QueryMessageProcessor,读取消息时主要根据 MessageId 中解析出的 Broker 地址和 CommitLog Offset 定位消息记录,并解析后返回。

6.2 按照Message Key查询消息

“按照Message Key查询消息”,主要是基于RocketMQ的IndexFile索引文件来实现的。RocketMQ的索引文件逻辑结构,类似JDK中HashMap的实现。索引文件的具体结构如下:

!

IndexFile 为用户提供按照 Message Key 查询消息的索引能力,但由于底层使用 Hash 索引,查询结果仍需读取消息并校验真实 Key,IndexFile文件的存储位置是:$HOME\store\index${fileName},文件名fileName是以创建时的时间戳命名的,文件大小是固定的,等于40+500W*4+2000W*20= 420000040个字节大小。如果消息的properties中设置了UNIQ_KEY这个属性,就用 topic + “#” + UNIQ_KEY的value作为 key 来做写入操作。如果消息设置了KEYS属性(多个KEY以空格分隔),也会用 topic + “#” + KEY 来做索引。

其中的索引数据包含了Key Hash/CommitLog Offset/Timestamp/NextIndex offset 这四个字段,一共20 Byte。NextIndex offset 即前面读出来的 slotValue,如果有 hash冲突,就可以用这个字段将所有冲突的索引用链表的方式串起来了。Timestamp记录的是消息storeTimestamp之间的差,并不是一个绝对的时间。整个Index File的结构如图,40 Byte 的Header用于保存一些总的统计信息,4*500W的 Slot Table并不保存真正的索引数据,而是保存每个槽位对应的单向链表的头。20*2000W 是真正的索引数据,即一个 Index File 可以保存 2000W个索引。

“按照Message Key查询消息”的方式,RocketMQ的具体做法是,主要通过Broker端的QueryMessageProcessor业务处理器来查询,读取消息的过程就是用topic和key找到IndexFile索引文件中的一条记录,根据其中的commitLog offset从CommitLog文件中读取消息的实体内容。

四、消息从生产到消费的完整过程

一条普通消息的完整链路如下:

1. Producer 启动
2. Producer 从 NameServer 获取 Topic 路由
3. Producer 选择一个 MessageQueue
4. Producer 将消息发送给 Queue 所在 Broker
5. Broker 把消息追加到 CommitLog
6. Broker 建立 ConsumeQueue 和索引
7. Broker 返回发送结果
8. Consumer 从 Broker 获取消息
9. Consumer 执行业务逻辑
10. Consumer 提交消费结果
11. Broker 更新该 Consumer Group 的消费进度

4.1 Producer 如何选择 Queue

普通消息通常通过客户端负载均衡选择 Queue。假设一个 Topic 在两个 Broker 上共有八个 Queue,Producer 会在这些 Queue 中分散发送,从而提高吞吐量。

如果是顺序消息,则不能随机选择,而应该根据业务键稳定选择 Queue:

queueIndex = hash(orderId) % queueCount;

只要同一个 orderId 的 Hash 结果稳定,它的消息就会进入同一个 Queue。

需要注意,直接使用 % queueCount 时,如果 Queue 数量变化,映射结果也可能改变。因此企业系统调整 Queue 数量时,要评估顺序业务在扩缩容期间的影响。

4.2 PushConsumer 本质上并不是 Broker 主动长连接推送

传统 RocketMQ 中常见 DefaultMQPushConsumer 或 Spring 的监听器,从业务视角看像“消息主动推送到方法”,但客户端底层仍然包含拉取、长轮询、缓存和消费线程池等机制。

可以理解为:

客户端向 Broker 发起拉取请求
      ↓
没有消息时 Broker 暂时挂起请求
      ↓
新消息到达或超时
      ↓
Broker 返回消息
      ↓
客户端放入本地缓存
      ↓
消费线程池调用监听器

RocketMQ 5.x 官方文档将消费者分为 PushConsumer、SimpleConsumer 和 PullConsumer,不同消费者类型在获取方式、消费确认和重试机制上有所不同。

在普通 Spring Boot 业务项目中,使用监听器式 PushConsumer 最方便;需要精确控制拉取、批量处理、消费确认或流量时,可以考虑 SimpleConsumer 或 PullConsumer。


五、RocketMQ 为什么性能高:CommitLog、ConsumeQueue 与顺序写盘

RocketMQ 的高吞吐不是因为“不写磁盘”,而是因为它尽量把随机写变成顺序写,并利用操作系统 Page Cache、内存映射文件和批量刷盘提高性能。

5.1 CommitLog:所有消息的物理存储文件

Broker 接收到的消息会按照到达顺序追加写入 CommitLog:

CommitLog
┌──────────────────────────────────────────┐
│ TopicA-Q0-message1                       │
│ TopicB-Q2-message1                       │
│ TopicA-Q1-message2                       │
│ TopicC-Q0-message1                       │
│ TopicA-Q0-message3                       │
└──────────────────────────────────────────┘

不同 Topic、不同 Queue 的消息并不是分别保存成完整消息文件,而是统一顺序追加到 CommitLog。这样能够显著减少磁盘随机写。

每条 CommitLog 记录通常包含:

消息总长度
消息 ID
Topic
Queue ID
Queue Offset
消息体
Tag
Key
属性
时间戳
物理偏移量
CRC 等校验信息

5.2 ConsumeQueue:逻辑消费索引

如果所有消费者每次都扫描整个 CommitLog,效率会很低。因此 RocketMQ 为每个 Topic 的每个 Queue 建立 ConsumeQueue:

consumequeue/
└── order-event-topic/
    ├── 0/
    ├── 1/
    ├── 2/
    └── 3/

ConsumeQueue 中保存的不是完整消息,而是轻量索引,核心指向:

消息在 CommitLog 中的物理偏移
消息长度
Tag Hash 等辅助信息

消费者拉取 order-event-topic 的 Queue 0 时,先顺序读取 Queue 0 对应的 ConsumeQueue,再根据物理偏移去 CommitLog 读取完整消息。

所以可以记成:

CommitLog:真正存放完整消息
ConsumeQueue:按 Topic + Queue 组织的消费索引

5.3 IndexFile:按 Key 查询消息

IndexFile(哈希索引文件)仅用于根据消息 KEY 精确检索消息,普通消费完全不用它,只有运维控制台、API 按 key 查消息才走这个文件 关键前提:只有你代码里设置了 RocketMQHeaders.KEYS=eventId,Broker 才会给这条消息生成 Index 索引;消息不带 KEY,不会写入 IndexFile,后台就搜不到这条消息。

IndexFile 物理文件基础信息

  1. 存放路径:store/index/时间戳命名文件
  2. 单文件固定大小 ≈ 400MB,写满自动新建下一个;文件名是创建时间戳,用于快速按时间筛选索引文件。
  3. 底层结构 = 磁盘版 HashMap + 单向链表解决哈希冲突,头插法写入新索引,新消息链表在前,查找更快 整体三段:IndexHeader(40B) + HashSlots(500万×4B) + Index条目区(最多2000万×20B)

1)IndexHeader 文件头(固定 40 字节,元数据)

字段按字节拆分:

  • 8B beginTimestamp:本文件第一条消息存储时间
  • 8B endTimestamp:本文件最后一条消息存储时间
  • 8B beginPhyoffset:第一条消息在 CommitLog 的物理偏移量
  • 8B endPhyoffset:最后一条消息在 CommitLog 偏移量
  • 4B hashSlotCount:被占用过的哈希槽数量
  • 4B indexCount:当前文件总索引条目数

作用:快速判断当前 IndexFile 的时间区间、CommitLog 偏移范围,查询时跳过不匹配的索引文件。

2)HashSlots 哈希槽区(固定 500 万个槽,每个 4 字节)

  • 总大小:5000000 × 4B = 20MB
  • 每个 slot 只存4 字节整数对应链表最新一条 Index 条目的序号,也就是条目中的位置
  • 哈希计算规则: eventId.hashCode() 是 Java 自带方法:把字符串 eventId 转换成一个整数哈希值。 举例子: eventId = “uuid-001” → 算出一个数字 123456 eventId = “uuid-002” → 算出另一个数字 789012 key.hashCode() % 5000000 = 对应slot下标
  • 冲突处理:同 hash 值的 key 共用一个 slot,slot 存链表头节点索引,旧索引通过preIndexNo串联成单向链表

3)Index 条目区(每条 20 字节,一条 KEY 对应一条条目) 只会往后追加,不会覆盖、不会中间插入,物理上完全有序、连续,和普通数组一模一样。 单条 Index 固定 20B 结构:

  1. 4B keyHash:业务 KEY(你的 eventId)的哈希值
  2. 8B phyOffset:消息完整内容在 CommitLog 里的全局物理偏移量(核心,用来找到原始消息)
  3. 4B timeDiff:消息存储时间 - 当前 IndexFile 的 beginTimestamp(节省存储空间)
  4. 4B preIndexNo:前一条冲突索引的序号,形成冲突链表

5.3.1 完整写入流程

  1. 生产者发送消息,设置 RocketMQHeaders.KEYS = event.getEventId()
  2. Broker 把整条消息完整写入 CommitLog,拿到这条消息的全局物理偏移量 phyOffset
  3. 异步后台线程(ReputMessageService) 扫描新写入 CommitLog 的消息:
    • 如果消息 Header 存在 KEY(eventId),进入索引构建;无 KEY 直接跳过;
  4. 计算 KEY 哈希:hash = eventId.hashCode()
  5. 对 500 万取余,定位到对应 HashSlot;
  6. 在 Index 条目区末尾追加一条 20B 索引,记录 keyHash、phyOffset、timeDiff、preIndexNo
  7. 更新对应 Slot 的值为当前新条目的序号(头插法,新索引变成链表头部);
  8. Mmap 内存映射,定时刷盘持久化到 IndexFile 磁盘文件。

重点:IndexFile 是异步构建,写入 CommitLog 和生成索引不是同步,刚发消息立刻查 KEY 可能短暂查不到。

5.3.2 完整写入流程根据 KEY 查询消息完整流程

需求:输入 eventId,查出完整 FileUploadedEvent 消息

  1. 传入两个条件:KEY=eventId + 时间范围(避免遍历全部 IndexFile);

  2. 根据时间区间筛选匹配的 IndexFile 文件(文件名是时间戳,快速过滤);

  3. 对每个匹配的 IndexFile: ① 计算 KEY 哈希 hash = eventId.hashCode(),取模 500 万得到 slot 位置; ② 读取 slot 里存储的 index 条目序号(链表头); ③ 循环读取该序号对应的 Index 条目:

    • 对比条目内keyHash和查询 KEY 的 hash 是否一致;
    • 一致则取出phyOffset,去 CommitLog 读取完整消息;
    • 不一致通过preIndexNo找链表上一条索引,继续对比; ④ 遍历完链表无匹配则换下一个 IndexFile;
  4. 通过 phyOffset 读取 CommitLog 二进制,反序列化得到完整 event 对象返回给查询者。

5.4 内存映射和 Page Cache

RocketMQ 通常不会让 Java 业务线程每次都直接执行传统磁盘随机 I/O,而会借助内存映射文件把文件映射到进程虚拟内存,再通过操作系统 Page Cache 管理实际内存页。 简化理解:

应用写入映射内存
      ↓
数据进入操作系统 Page Cache
      ↓
操作系统或刷盘线程写入磁盘

这样能够减少用户态与内核态之间的数据复制和系统调用开销。不过 Page Cache 不是永久存储,机器突然断电时,尚未真正落盘的数据仍有风险,因此 RocketMQ 又提供同步刷盘和异步刷盘策略。

5.5 同步刷盘与异步刷盘

异步刷盘的流程大致是:

消息写入 Page Cache
      ↓
Broker 返回发送成功
      ↓
后台线程批量刷入磁盘

优点是吞吐量高、延迟低;风险是操作系统还没完成刷盘时机器突然故障,极少量消息可能丢失。

同步刷盘的流程大致是:

消息写入 Page Cache
      ↓
等待消息真正刷入磁盘
      ↓
Broker 返回发送成功

优点是可靠性更高,缺点是发送延迟增加、吞吐量下降。官方配置说明中将 Broker 刷盘方式区分为 SYNC_FLUSHASYNC_FLUSH

实际选型应根据业务决定:

通知、日志、非核心事件:
    通常更关注吞吐量,可以接受异步刷盘

支付、结算、资金相关事件:
    更关注可靠性,需要结合同步刷盘、同步复制和业务补偿

不能简单认为开启同步刷盘就“绝对不丢消息”,因为可靠性还取决于 Producer 是否发送成功、主从复制、Broker 故障切换、消费确认和业务幂等。


六、消息发送方式与 Spring Boot 使用

传统 Spring Boot 项目通常使用 Apache RocketMQ Spring 集成:

<dependency>
    <groupId>org.apache.rocketmq</groupId>
    <artifactId>rocketmq-spring-boot-starter</artifactId>
    <version>${rocketmq-spring.version}</version>
</dependency>

官方项目说明该 Starter 支持同步、异步、单向、顺序、批量、事务和延时消息等能力。版本应根据当前 Spring Boot、RocketMQ Client 和服务端兼容关系统一管理,不建议复制旧教程中的固定版本。

[!note] 一条 RocketMQ 消息底层长什么样(简化结构)

【消息元数据(Header)】
KEYS = "uuid-xxx-eventId"
fileId = "file_10086"
schemaVersion = 1
TAG = "FILE_UPLOADED"
MSG_ID = "Broker生成的全局消息ID"
投递时间、重试次数、事务标记等系统内置头

【消息体 Body(二进制)】
{
  "eventId": "uuid-xxx-eventId",
  "fileId": "file_10086",
  "fileName": "a.pdf",
  "fileSize": 102400,
  "uploadTime": "2026-07-14"
}

  • Header:轻量键值对,读取速度快,不需要反序列化整个对象
  • Body:序列化后的完整业务对象,量大,处理需要 JSON / 反序列化

6.0搭建rocketmq-producer(消息生产者)

修改配置文件application.yml

spring:
  application:
    name: rocketmq-producer
rocketmq:
  # rocketMq的nameServer地址
  name-server: 127.0.0.1:9876     
  producer:
    # 生产者组别
    group: test-group
    # 消息发送的超时时间
    send-message-timeout: 3000
    # 异步消息发送失败重试次数
    retry-times-when-send-async-failed: 2
    # 发送消息的最大大小,单位字节,这里等于4M
    max-message-size: 4194304      

生产环境不应在代码中写死 NameServer 地址,可以通过 Nacos 配置、环境变量或部署平台注入:

rocketmq:
  name-server: ${ROCKETMQ_NAME_SERVER}

在测试类里面测试发送消息 往test主题里面发送一个简单的字符串消息

/**
 * 注入rocketMQTemplate,我们使用它来操作mq
 */
@Autowired
private RocketMQTemplate rocketMQTemplate;
 
/**
 * 测试发送简单的消息
 *
 * @throws Exception
 */
@Test
public void testSimpleMsg() throws Exception {
    // 往test的主题里面发送一个简单的字符串消息
    // syncSend同步消息
    // asyncSend异步
    // 参数1:topic名 参数2:Object,数据
    SendResult sendResult = rocketMQTemplate.syncSend("test", "我是一个简单的消息");
    // 拿到消息的发送状态
    System.out.println(sendResult.getSendStatus());
    // 拿到消息的id
    System.out.println(sendResult.getMsgId());
}

运行后查看控制台

6.0搭建rocketmq-consumer(消息消费者)

修改配置文件application.yml

spring:
  application:
    name: rocketmq-consumer
rocketmq:
  name-server: 127.0.0.1:9876
#    consumer:
#        group: aaa-group 不需要写,一个项目一般有很多消费者组

监听器SimpleMsgListener 消费者要消费消息,就添加一个监听器,SpringBoot一启动,监听器就开始持续工作 1、类上添加注解 @Component 和 @RocketMQMessageListener 。topic指定消费的主题,consumerGroup 指定消费组,一个主题可以有多个消费者组,一个消息可以被多个不同的组的消费者都消费 2、实现 RocketMQListener 接口,泛型可以为具体的数据类型,如果想拿到消息的其他参数(如消息头、消息体,例如key之类的),泛型用MessageExt

package com.test.listener;
 
import org.apache.rocketmq.spring.annotation.MessageModel;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;
 
@Component
@RocketMQMessageListener(topic = "test", consumerGroup = "test-group", messageModel = MessageModel.CLUSTERING)
public class SimpleMsgListener implements RocketMQListener<String> {
 
    /**
     * 消费消息的方法
     * 
     * @param message 消息内容,类型和上面的泛型一致。如果泛型指定了固定的类型,消息体就是我们的参数
     */
    @Override
    public void onMessage(String message) {
        System.out.println(message);
    }
}

泛型使用MessageExt

@Component
@RocketMQMessageListener(topic = "test", consumerGroup = "test-group", messageModel = MessageModel.CLUSTERING)
public class SimpleMsgListener implements RocketMQListener<MessageExt> {
 
    /**
     * 消费消息的方法
     * 
     * @param message 消息内容,类型和上面的泛型一致。如果泛型指定了固定的类型,消息体就是我们的参数
     */
    @Override
    public void onMessage(MessageExt message) {
        System.out.println(new String(message.getBody()));
    }
}

onMessage方法没有返回值,要表示消息签收?

  • 方法报错就拒收
  • 方法不报错就签收

6.1 定义统一事件对象

不要直接发送数据库 Entity,也不要随意发送没有版本信息的 Map。推荐定义稳定的事件 DTO:

方式 示例 推荐度 原因
直接发送 Entity send(entity) ❌ 强烈不推荐 太危险,紧耦合,容易泄露敏感数据
发送 Map send(map) ❌ 不建议 无类型安全,无版本控制,维护困难
发送事件 DTO send(FileUploadedEvent) ✅ 强烈推荐 稳定、自描述、可演进、解耦
事件 DTO 是生产者和消费者之间的"契约",它稳定、自描述、可演进
@Data
@Builder
@NoArgsConstructor
@AllArgsConstructor
public class FileUploadedEvent {

    /**
     * 全局唯一事件 ID,用于幂等和排障。
     */
    private String eventId;

    /**
     * 文件业务 ID。
     */
    private Long fileId;

    private Long userId;

    private String fileName;

    /**
     * 事件结构版本,便于以后兼容升级。
     */
    private Integer schemaVersion;

    private LocalDateTime occurredAt;
}

企业事件对象应尽量具有:

eventId
businessId
eventType
schemaVersion
occurredAt
traceId
业务数据

事件应该表达已经发生的业务事实,例如 FileUploadedEventOrderPaidEvent,而不是模糊的 FileMessage

6.2 同步发送

同步发送会等待 Broker 返回结果,适合需要明确知道消息是否发送成功的业务: 把事件唯一 ID 设置为 RocketMQ 消息的 Key,主要用于查询、排查和追踪消息

import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.spring.core.RocketMQTemplate;
import org.apache.rocketmq.spring.support.RocketMQHeaders;
import org.springframework.messaging.Message;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.stereotype.Service;

/**
 * 文件事件生产者
 * 负责发送文件相关事件到RocketMQ消息队列
 * Topic:file-event-topic,TAG:FILE_UPLOADED 区分消息业务类型
 */
@Service                    // ① 将当前类注册为Spring容器Bean,可被@Autowired注入使用
@RequiredArgsConstructor    // ② Lombok注解:生成仅包含final修饰字段的构造器,实现构造器注入(推荐,优于@Autowired字段注入)
@Slf4j                       // ③ Lombok注解:自动生成log日志对象,无需手动声明Logger
public class FileEventProducer {

    // 常量定义Topic名称,统一管理,避免硬编码散落在代码各处
    private static final String TOPIC = "file-event-topic";
    // 消息TAG:用于同一Topic下过滤不同业务消息(文件上传事件标记)
    private static final String TAG_FILE_UPLOAD = "FILE_UPLOADED";

    // RocketMQ操作模板类,SpringRocketMQ自动装配,封装发送、事务、批量等API
    private final RocketMQTemplate rocketMQTemplate;

    /**
     * 发送【文件已上传】事件消息
     * @param event 文件上传事件实体,承载业务数据
     */
    public void sendFileUploaded(FileUploadedEvent event) {
        // 1. 组装消息目的地:格式规则 Topic:TAG,RocketMQ标准分隔符
        String destination = TOPIC + ":" + TAG_FILE_UPLOAD;

        // 2. 构建Spring标准Message消息对象
        // withPayload(event):设置消息体,真正要投递的业务实体FileUploadedEvent
        Message<FileUploadedEvent> message = MessageBuilder.withPayload(event)
                // RocketMQ消息KEY:消息唯一业务标识,用于消息查询、日志追踪、死信定位
                把事件唯一 ID 设置为 RocketMQ 消息的 Key,主要用于查询、排查和追踪消息。
                两个参数一个,键值对,一个是键,一个是值
                .setHeader(RocketMQHeaders.KEYS, event.getEventId())
                // 自定义消息头1:文件ID,消费者可直接从Header快速获取,无需解析消息体
                .setHeader("fileId", event.getFileId())
                // 自定义消息头2:消息数据结构版本,用于后续消息结构兼容、灰度升级
                .setHeader("schemaVersion", 1)
                // 完成消息构建,返回Message实例
                .build();

        // 3. 同步发送消息 syncSend:阻塞等待Broker返回发送结果,可靠性最高
        SendResult result = rocketMQTemplate.syncSend(destination, message);

        // 4. 打印发送成功日志,埋点关键业务字段,方便线上排查问题
        log.info(
                "文件上传事件发送成功, eventId={}, fileId={}, msgId={}, sendStatus={}",
                event.getEventId(),       // 业务事件唯一ID(自定义业务主键)
                event.getFileId(),        // 文件业务ID
                result.getMsgId(),        // RocketMQ服务端生成的全局唯一消息ID
                result.getSendStatus()    // 发送状态:SEND_OK 代表投递成功
        );
    }
}

同步发送成功表示 Broker 已按当前配置接受消息,不代表消费者已经完成业务。

发送方必须处理异常:

try {
    rocketMQTemplate.syncSend(destination, message);
} catch (Exception e) {
    log.error(
            "消息发送失败, eventId={}, destination={}",
            event.getEventId(),
            destination,
            e
    );

    throw new MessageSendException(
            "发送文件事件失败",
            e
    );
}

不要捕获异常后只打印日志并继续返回业务成功,否则可能出现数据库已经提交、消息却没有发送的情况。

6.3 异步发送

异步发送不会长时间阻塞当前线程,通过回调获取结果:

// 发送异步消息,发送完以后会有一个异步通知
rocketMQTemplate.asyncSend(
        "notification-topic:EMAIL", // ① 目的地topic
        event,  // ② 消息内容
        new SendCallback() {    // ③ 回调对象(重点!)
			// 成功时的处理
            @Override
            public void onSuccess(SendResult sendResult) {
                log.info(
                        "异步消息发送成功, eventId={}, msgId={}",
                        event.getEventId(),
                        sendResult.getMsgId()
                );
            }
		 // 失败时的处理
            @Override
            public void onException(Throwable throwable) {
                log.error(
                        "异步消息发送失败, eventId={}",
                        event.getEventId(),
                        throwable
                );
            }
        }
);

异步发送适合主线程不能被 Broker 网络响应阻塞,但应用必须认真处理失败回调。异步不等于不管结果。

6.4 单向发送

单向发送只负责把请求写出,不等待 Broker 响应:

  • 吞吐量很大,存在消息丢失的风险,可用于日志信息的发送
// 发送单向消息,没有返回值和结果
rocketMQTemplate.sendOneWay(
        "log-topic:ACCESS_LOG",
        accessLogEvent
);

它延迟最低,但生产者无法确认发送结果。适合允许少量丢失的非核心日志、监控或统计数据,不适合支付、订单和文件元数据等核心业务。

6.5 发送消息时设置 Key、Header 与 Tag

Message<FileUploadedEvent> message =
        MessageBuilder.withPayload(event)
                .setHeader(
                        RocketMQHeaders.KEYS,
                        event.getEventId()
                )
                .setHeader("fileId", event.getFileId())
                .setHeader("schemaVersion", 1)
                .build();

rocketMQTemplate.syncSend(
        "file-event-topic:FILE_UPLOADED",
        message
);

注意不要在 Header 中放大量数据,消息主体才是主要业务数据。AccessKey、密码、完整 Token 等敏感信息不能进入消息或日志。

6.6发送延迟消息

4.x版本只支持特定等级的延迟。延迟等级,从1级开始分别对应:1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h

/**
 * 测试延迟消息
 *
 * @throws Exception
 */
@Test
public void testDelay() throws Exception {
    // 构建消息对象
    Message<String> message = MessageBuilder.withPayload("我是一个延迟消息").build();
    // 发送一个延时消息,延迟等级为4级,也就是30s后被监听消费
    // 参数3:连接MQ的超时时间 参数4:延迟等级
    SendResult sendResult = rocketMQTemplate.syncSend("test", message, 2000, 4);
    System.out.println(sendResult.getSendStatus());
}

运行后,查看消费者端,过了30s才被消费

5.x版本支持任意时间的延迟

一段时间后消费消息 destination格式_:_ topicName:tags,没有tag的话,直接写topicName即可

 rocketMQTemplate.syncSendDelayTimeSeconds(
        "topic",
        message,
        10
);

6.7 发送事务消息

而事务对象(Transaction Object) 是另一个层面的概念,它负责管理事务的上下文和状态

6.7.1 事务消息解决什么核心问题

普通消息两大矛盾(无法保证原子性):

  1. 先发消息,再执行本地 DB 事务 消息发送成功,本地数据库事务报错回滚 → 下游收到无效消息,脏数据
  2. 先执行本地 DB 事务,再发消息 数据库操作成功,网络宕机消息发不出去 → 下游收不到通知,业务断流

事务消息目标:本地数据库事务 + MQ 消息发送 要么同时成功、要么同时失败,保证分布式最终一致性 底层方案:2PC 两阶段提交 + 事务回查补偿机制

[!note] 1.Half Message 半消息

生产者发给 Broker,但消费者完全看不见的消息。 Broker 不会存入你的业务 Topic,而是存在系统内部 Topic: RMQ_SYS_TRANS_HALF_TOPIC(半消息专用存储队列) 特点:消息持久化落地,但标记「待确认,禁止投递」。

[!note] 2.OP Topic(事务操作日志)

系统内部 Topic:RMQ_SYS_TRANS_OP_HALF_TOPIC 用来记录半消息最终状态:commit/rollback。 Broker 定时扫描长期没有 OP 记录的半消息,触发事务回查。

[!note] 3.三种事务状态(生产者回调返回)

  1. COMMIT_MESSAGE:本地事务成功 → Broker 把半消息转移到真实业务 Topic,消费者正常消费
  2. ROLLBACK_MESSAGE:本地事务失败 → Broker 直接删除半消息,下游永远收不到
  3. UNKNOWN:状态未知(宕机、超时、异常卡住)→ Broker 定时主动回查生产者确认状态

[!note] 4.事务回查(补偿核心)

生产者宕机、网络中断,导致 Broker 没收到 commit/rollback 指令; Broker 每隔固定时间(默认 1 分钟)主动调用生产者接口,查询这条消息对应的本地数据库事务是否执行成功。

6.7.2 完整 4 步标准流程(分正常流程、异常回查流程)

阶段 1:发送半消息(Prepare 准备阶段)

  1. 生产者调用 sendMessageInTransaction 发送事务消息
  2. Broker 收到消息,存入系统半消息队列 RMQ_SYS_TRANS_HALF_TOPIC,持久化 CommitLog
  3. Broker 返回 ACK 给生产者:半消息投递成功

关键:此时消费者完全收不到这条消息

阶段 2:执行本地数据库事务 生产者收到半消息成功 ACK 后,自动回调 executeLocalTransaction 方法执行本地业务 DB 事务: 示例:创建订单、扣减库存、更新支付记录(加 @Transactional 数据库事务) 执行完返回 3 种状态之一:COMMIT / ROLLBACK / UNKNOWN

阶段 3:上报最终确认(Commit/Rollback)

  1. 返回COMMIT:生产者给 Broker 发提交指令 Broker 操作:把半消息从内部半队列迁移到真实业务 Topic,消费者开始拉取消费
  2. 返回ROLLBACK:生产者给 Broker 发回滚指令 Broker 操作:直接丢弃这条半消息,下游永远看不到
  3. 返回UNKNOWN:不做任何确认,Broker 等待定时回查

阶段 4:事务回查补偿(宕机 / 网络丢失场景) 触发条件:Broker 长期没收到这条半消息的 commit/rollback 指令

  1. Broker 主动发起回查请求,调用生产者 checkLocalTransaction 回调
  2. 生产者根据消息唯一 KEY(eventId)查询本地数据库:
    • 查到订单记录 → 返回 COMMIT
    • 数据库无记录、事务回滚 → 返回 ROLLBACK
    • 业务还在处理 → 继续返回 UNKNOWN,等待下一轮回查
  3. Broker 根据回查结果处理半消息,解决生产者宕机丢失确认的问题

6.7.3 SpringBoot 代码标准实现

1.事务监听器核心(必须实现两个方法)

@Slf4j
@Component
// 绑定事务生产者组,一组生产者共用一个监听器
@RocketMQTransactionListener(txProducerGroup = "file-tx-producer-group")
public class FileTxListener implements RocketMQLocalTransactionListener {

    @Autowired
    private FileMapper fileMapper; // 本地数据库Mapper

    // 阶段2:半消息发送成功后执行本地事务
    @Override
    @Transactional(rollbackFor = Exception.class)
    public RocketMQLocalTransactionState executeLocalTransaction(Message msg, Object arg) {
        FileUploadedEvent event = (FileUploadedEvent) msg.getPayload();
        String eventId = (String) msg.getHeaders().get(RocketMQHeaders.KEYS);
        try {
            // 本地数据库事务:写入文件上传记录
            FileRecord record = new FileRecord();
            record.setEventId(eventId);
            record.setFileId(event.getFileId());
            fileMapper.insert(record);
            // 本地事务成功,通知Broker提交消息
            return RocketMQLocalTransactionState.COMMIT;
        } catch (Exception e) {
            log.error("本地文件事务失败 eventId:{}", eventId, e);
            // DB失败,通知Broker丢弃消息
            return RocketMQLocalTransactionState.ROLLBACK;
        }
    }

    // 阶段4:Broker回查时调用,校验本地事务是否存在
    @Override
    public RocketMQLocalTransactionState checkLocalTransaction(Message msg) {
        String eventId = (String) msg.getHeaders().get(RocketMQHeaders.KEYS);
        // 根据eventId查询数据库是否存在这条上传记录
        FileRecord record = fileMapper.selectByEventId(eventId);
        if (record != null) {
            // DB存在记录 → 本地事务成功,提交消息
            return RocketMQLocalTransactionState.COMMIT;
        } else {
            // DB无记录 → 本地事务回滚,丢弃消息
            return RocketMQLocalTransactionState.ROLLBACK;
        }
    }
}

  1. 生产者发送事务消息(替换你原来的 syncSend)
@Service
@RequiredArgsConstructor
@Slf4j
public class FileEventTxProducer {
    private static final String TOPIC = "file-event-topic";
    private final RocketMQTemplate rocketMQTemplate;

    public void sendTxFileUpload(FileUploadedEvent event) {
        String dest = TOPIC + ":FILE_UPLOADED";
        // 构建消息,eventId作为KEY存入索引
        Message<FileUploadedEvent> message = MessageBuilder.withPayload(event)
                .setHeader(RocketMQHeaders.KEYS, event.getEventId())
                .setHeader("fileId", event.getFileId())
                .setHeader("schemaVersion", 1)
                .build();
        // 发送事务消息,第二个参数为自定义业务参数(可传null)
        rocketMQTemplate.sendMessageInTransaction(dest, message, null);
    }
}

事务消息先用半消息把消息预存在 Broker 但不让下游消费;再同步执行本地数据库事务;根据 DB 结果决定消息是否对外投递;生产者宕机时依靠 Broker 定时回查数据库状态兜底,最终实现本地业务与 MQ 消息同时成功 / 同时失败。

七、消费模式、消费进度与负载均衡

7.1 集群消费

集群消费是企业中最常见的模式。同一个 Consumer Group 中的多个实例共同消费消息:

Queue 0 → Consumer A
Queue 1 → Consumer B
Queue 2 → Consumer A
Queue 3 → Consumer B

一条消息在一个消费组内通常只交给一个消费者实例处理。启动更多实例可以提高消费能力,但消费者实例数长期大于 Queue 数量时,多出的实例可能分不到 Queue。

Spring 消费者示例:

@Component
@Slf4j
@RocketMQMessageListener(
        topic = "file-event-topic",              // ① 订阅的主题
        selectorExpression = "FILE_UPLOADED",    // ② 过滤标签,只消费带 `FILE_UPLOADED` Tag 的消息。
        consumerGroup = "networkdisk-ai-file-consumer", // ③ 消费组,消费者组名称,用于管理消费者集群。
        messageModel = MessageModel.CLUSTERING,  // ④ 消费模式,集群消费,组内多个消费者**负载均衡**消费消息
        consumeMode = ConsumeMode.CONCURRENTLY   // ⑤ 消费并发模式,多线程并发消费,**默认推荐**,性能高
)
//实现 RocketMQ 的监听接口,指定反序列化类型。
public class FileUploadedConsumer
        implements RocketMQListener<FileUploadedEvent> {

    @Override
    public void onMessage(FileUploadedEvent event) {
        log.info(
                "开始处理文件上传事件, eventId={}, fileId={}",
                event.getEventId(),
                event.getFileId()
        );

        // 执行文本解析、向量化或其他异步任务
        processFile(event);

        log.info(
                "文件上传事件处理成功, eventId={}, fileId={}",
                event.getEventId(),
                event.getFileId()
        );
    }

    private void processFile(FileUploadedEvent event) {
        // 业务逻辑
    }
}
  1. RocketMQ 拉取到消息
  2. 反序列化成 FileUploadedEvent 对象
  3. 调用 onMessage(event) 方法
  4. 执行业务逻辑(processFile)
  5. 方法执行完毕 → RocketMQ 自动提交消费位点(Offset)
┌─────────────────────────────────────────────────────────────────┐
│                    RocketMQ 消费流程                           │
├─────────────────────────────────────────────────────────────────┤
│                                                                 │
│  Broker(RocketMQ 服务端)                                     │
│       │                                                        │
│       │ 1. 推送消息给消费者                                    │
│       ▼                                                        │
│  ┌─────────────────────────────────────────────────────────────┐│
│  │ FileUploadedConsumer(你的消费端)                         ││
│  │                                                             ││
│  │ 2. @RocketMQMessageListener 启动监听                       ││
│  │    ├── topic = "file-event-topic"                         ││
│  │    ├── selectorExpression = "FILE_UPLOADED"               ││
│  │    └── consumerGroup = "networkdisk-ai-file-consumer"     ││
│  │                                                             ││
│  │ 3. 拉取到消息(JSON 格式)                                 ││
│  │    {                                                       ││
│  │      "eventId": "evt-001",                                ││
│  │      "fileId": 12345,                                     ││
│  │      "fileName": "report.pdf"                             ││
│  │    }                                                       ││
│  │                                                             ││
│  │ 4. 自动反序列化成 FileUploadedEvent 对象                  ││
│  │                                                             ││
│  │ 5. onMessage(event) 被调用                                ││
│  │    ├── 日志:开始处理                                     ││
│  │    ├── processFile(event)                                 ││
│  │    │   ├── 下载文件                                       ││
│  │    │   ├── 解析文本                                       ││
│  │    │   ├── 向量化                                         ││
│  │    │   └── 存入向量数据库                                 ││
│  │    └── 日志:处理成功                                     ││
│  │                                                             ││
│  │ 6. 方法正常返回 → 提交消费位点(Offset)                  ││
│  │                                                             ││
│  └─────────────────────────────────────────────────────────────┘│
│                                                                 │
└─────────────────────────────────────────────────────────────────┘

7.2 广播消费

广播消费表示同一个 Consumer Group 中的每个实例都处理每条消息:

Message 1 → Consumer A
          → Consumer B
          → Consumer C

适合刷新各节点本地缓存、更新本地配置、触发每台机器执行本地任务等场景。

messageModel = MessageModel.BROADCASTING

广播模式不能用来简单提高业务可靠性,因为每个实例都执行一次,容易造成重复扣款、重复发短信或重复修改数据库。

7.3 消费进度 Offset

RocketMQ 需要记录每个 Consumer Group 在每个 Queue 消费到什么位置:

ConsumerGroup: networkdisk-ai-file-consumer

Queue 0 offset = 1024
Queue 1 offset = 986
Queue 2 offset = 1102
Queue 3 offset = 1007

Offset 表示消费位置,而不是消息是否完成业务的唯一凭证。消费者成功确认后,消费进度才会向前推进。 不同 Consumer Group 有各自独立的消费进度,因此库存服务消费到第 1000 条,不影响通知服务从第 500 条继续消费。

7.4 消费负载均衡与 Rebalance

当消费者实例数量变化、Queue 数量变化或 Broker 路由变化时,RocketMQ 会重新分配 Queue,这个过程称为 Rebalance。

例如:

原来:
Consumer A → Queue 0、1、2、3

新增 Consumer B 后:
Consumer A → Queue 0、1
Consumer B → Queue 2、3

Rebalance 期间可能发生短暂暂停、Queue 转移或重复投递。因此消费者必须是幂等的,不能假设一条消息在任何异常情况下绝对只会执行一次。


八、消息可靠性:如何处理消息丢失、重复和积压

企业面试中经常问:“RocketMQ 怎么保证消息不丢?”正确答案不能只说“持久化”,而应该分生产、存储和消费三个阶段分析。

Producer 发送阶段
        ↓
Broker 存储阶段
        ↓
Consumer 消费阶段

8.1 生产阶段可靠性

生产阶段可能出现:

业务数据库提交成功,但消息没发出去
消息已到 Broker,但 Producer 响应超时
Producer 进程在发送前宕机
网络闪断
发送失败重试后仍然失败

RocketMQ 客户端可以自动进行有限次数的发送重试,但官方建议重试次数不要无限增大,以免长时间阻塞业务线程;超过最大重试后,业务侧仍应进行回查和补偿。

核心业务不能只依赖:

saveOrder();
rocketMQTemplate.syncSend(...);

因为数据库事务和 RocketMQ 发送不是同一个本地事务。可能发生:

订单提交成功
      ↓
应用宕机
      ↓
消息没有发送

常见企业解决方案有三种。

方案一:本地消息表

在同一个数据库事务中同时保存业务数据和待发送事件:

本地事务:
    INSERT order
    INSERT outbox_event
提交成功

后台任务不断扫描未发送事件,发送到 RocketMQ;发送成功后更新状态。

outbox_event

id
event_id
topic
tag
payload
status
retry_count
next_retry_time
created_at

这种方案本质上把“数据库操作和记录发送意图”放在同一个本地事务里,然后通过重试保证最终发送成功。

方案二:事务消息

使用 RocketMQ 事务消息协调本地事务和消息提交,后面单独说明。

方案三:业务回查与对账补偿

对资金、订单等核心业务,定时检查业务状态与消息处理状态。例如订单已经支付但积分未增加时,重新生成补偿消息。 可靠系统不是依靠一个配置实现,而是:

发送确认
+ 有限重试
+ 本地记录
+ 定时补偿
+ 监控告警

8.2 Broker 存储阶段可靠性

Broker 存储可靠性取决于:

同步刷盘还是异步刷盘
主从同步复制还是异步复制
Broker 部署副本数
磁盘健康状况
故障切换能力

同步刷盘解决“消息是否写入当前节点磁盘”,同步复制解决“消息是否复制到从节点”。两者解决的问题不同。

同步刷盘:
    防止当前 Broker 进程或系统异常导致 Page Cache 数据丢失

同步复制:
    防止 Master 整机或磁盘损坏后消息没有副本

RocketMQ 还支持基于 Controller 的主从自动切换模式;官方文档说明,即使独立部署 Controller,仍需单独部署 NameServer 提供路由发现。

生产环境不能只部署单 Broker,否则 Broker 宕机后,消息发送和消费都会中断。

8.3 消费阶段可靠性

消费阶段最危险的错误是先确认成功,再执行业务:

错误顺序:
Consumer 确认消息成功
      ↓
执行业务
      ↓
业务失败

这样 Broker 认为消息已处理,不再重试,但实际业务没有完成。

正确顺序是:

执行业务
      ↓
业务成功
      ↓
返回消费成功

在 Spring 监听器中,如果业务失败,应抛出异常,让框架通知 Broker 本次消费失败:

@Override
public void onMessage(FileUploadedEvent event) {
    try {
        fileIndexService.buildIndex(event.getFileId());
    } catch (Exception e) {
        log.error(
                "构建文件索引失败, eventId={}, fileId={}",
                event.getEventId(),
                event.getFileId(),
                e
        );

        // 不要吞掉异常,否则框架可能认为消费成功
        throw e;
    }
}

下面这种写法非常危险:

@Override
public void onMessage(FileUploadedEvent event) {
    try {
        fileIndexService.buildIndex(event.getFileId());
    } catch (Exception e) {
        log.error("处理失败", e);

        // 错误:异常被吞掉,方法正常结束
    }
}

九、为什么消息一定要做幂等

RocketMQ 的可靠投递通常属于“至少一次”语义:消息不能轻易丢失,但可能被重复投递。

重复可能发生在:

Producer 发送成功,但响应超时,Producer 再次发送
Consumer 处理成功,但确认消息前宕机
网络异常导致 ACK 丢失
消费超时
Rebalance
人工重置 Offset
重试消息重新投递

因此消费者不能依赖“这条消息只执行一次”,必须通过业务设计保证重复执行不会产生错误结果。

9.1 数据库唯一键幂等

eventId 建唯一索引:

CREATE TABLE mq_consume_record (
    id BIGINT PRIMARY KEY AUTO_INCREMENT,
    event_id VARCHAR(64) NOT NULL,
    consumer_group VARCHAR(128) NOT NULL,
    consumed_at DATETIME NOT NULL,
    UNIQUE KEY uk_event_consumer (
        event_id,
        consumer_group
    )
);

消费时在本地事务中先插入消费记录,再执行业务:

@Transactional
public void consume(FileUploadedEvent event) {
    boolean inserted = consumeRecordMapper.insertIgnore(
            event.getEventId(),
            "networkdisk-ai-file-consumer"
    );

    if (!inserted) {
        // 已消费过,直接返回成功
        return;
    }

    fileIndexService.buildIndex(event.getFileId());
}

如果业务失败,事务回滚,消费记录也回滚;Broker 重试时仍可以重新执行。

9.2 业务状态机幂等

订单状态变化应使用条件更新:

UPDATE orders
SET status = 'PAID'
WHERE id = #{orderId}
  AND status = 'CREATED';

第一次消费更新成功,第二次重复消费因为状态已经是 PAID,影响行数为 0,不会重复执行状态迁移。

9.3 Redis 幂等的注意事项

可以使用 Redis SETNX

SET mq:consume:{consumerGroup}:{eventId} 1 NX EX 86400

但如果先写 Redis 成功,后续数据库操作失败,而 Redis Key 没有删除,下一次重试可能被误判为已经消费。因此核心业务更推荐让幂等记录与业务操作位于同一个数据库事务中,Redis 更适合非核心或允许短期去重的场景。

9.4 “Exactly Once”不能只靠 MQ 实现

消息系统可以减少重复投递,但要做到最终业务效果只发生一次,通常必须由:

MQ 投递语义
+ 业务唯一标识
+ 数据库唯一约束
+ 状态机
+ 本地事务

共同完成。

面试中可以回答:

RocketMQ 主要保证消息可靠投递,但在网络超时、重试和消费者故障场景下可能重复,因此业务层通常按照至少一次投递设计,通过事件 ID、唯一索引、状态机和本地事务实现最终业务幂等。


十、消费失败、重试消息与死信队列

10.1生产者重试

// 失败的情况重发3次(同步)
producer.setRetryTimesWhenSendFailed(3);
// 失败的情况重发3次(异步)
producer.setRetryTimesWhenSendAsyncFailed(3);
// 消息在1S内没有发送成功,就会重试
producer.send(msg, 1000);

【示例代码】

@Test
public void retryProducer() throws Exception {
    DefaultMQProducer producer = new DefaultMQProducer("retry-producer-group");
    producer.setNamesrvAddr(MqConstant.NAME_SRV_ADDR);
    producer.start();
    // 如果发送失败要重试几次(同步),不设置默认值是2
    producer.setRetryTimesWhenSendFailed(3);
    // 如果发送失败要重试几次(异步)
//        producer.setRetryTimesWhenSendAsyncFailed(3);
    String key = UUID.randomUUID().toString();
    System.out.println(key);
    Message message = new Message("retryTopic", "vip1", key, "我是vip666的文章".getBytes());
    producer.send(message);
    System.out.println("发送成功");
    producer.shutdown();
}

消费者处理失败后,RocketMQ 会按策略重新投递。重试不是立即无限循环,而是通过递增间隔给数据库、第三方服务或网络故障恢复时间。官方文档将消费失败后的重新投递定义为消费重试机制。

如果消息消费失败,默认会重试16次,重试的时间间隔:10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h

能否自定义重试次数?

可以,重试的次数一般设置为5

// 消费失败,重试几次
consumer.setMaxReconsumeTimes(5);

如果重试了16次(并发模式是16次,顺序模式下重试次数是 int 类型最大值) 都是失败的,怎么处理?

认为该消息是死信消息,将消息放在一个死信主题中去,名称:%DLQ%消费者组名,最后再实现一个消费者去消费死信消息,一般是发邮件发短信通知人工处理、做一些记录

死信队列只有一个队列

当消息处理失败的时候 该如何正确的处理?

方案一:处理死信队列,如果每个死信队列都写一个消费者,很麻烦

  • 由 RocketMQ 控制最大重试次数,达到阈值自动转死信;
  • 业务消费代码干净,不用判断重试次数;
  • 需要额外独立消费者程序监听死信 Topic;
  • 优点:解耦,死信单独运维,不阻塞正常业务消费;生产标准推荐方案。

/**
 * 方案一:独立死信队列消费者
 * 原理:消息重试耗尽阈值后自动进入 %DLQ%xxx 死信Topic;单独起消费者监听死信Topic统一兜底
 * 适用:业务不拦截重试,依赖RocketMQ自动转死信,单独程序人工巡检失败消息
 * @throws Exception
 */
@Test
public void retryDeadConsumer() throws Exception {
    // 创建Push消费者,分组:retry-dead-consumer-group(死信专用分组,与业务消费组隔离)
    DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("retry-dead-consumer-group");
    // 指定NameServer地址,连接MQ服务
    consumer.setNamesrvAddr(MqConstant.NAME_SRV_ADDR);
    // 订阅死信Topic:固定前缀%DLQ% + 原业务消费组名,*代表接收所有TAG
    consumer.subscribe("%DLQ%retry-consumer-group", "*");
    // 注册并发消息监听器
    consumer.registerMessageListener(new MessageListenerConcurrently() {
        @Override
        public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
            // 批量消息取第一条(测试Demo简化写法,生产需循环遍历所有msg)
            MessageExt messageExt = msgs.get(0);
            System.out.println("死信消息接收时间:" + new Date());
            // 打印原始消息体
            System.out.println("消息内容:" + new String(messageExt.getBody()));
            // 死信处理逻辑:落地数据库/文件记录、推送告警通知运维人工排查
            System.out.println("记录到特别的位置 文件 mysql 通知人工处理");
            // 返回成功:这条死信处理完毕,Broker删除,不再重复投递
            return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
        }
    });
    // 启动消费者
    consumer.start();
    // 阻塞主线程,防止测试方法直接结束
    System.in.read();
}

方案二:在实际生产过程中,一般重试3-5次,如果还没有消费成功,则可以把消息签收了,通知人工等处理

  • 完全代码控制重试阈值,不依赖 MQ 自动转死信;
  • 无需额外死信消费者,超限直接签收丢弃;
  • 优点:简单轻量化,小型项目快速实现;
  • 缺点:失败消息无独立队列沉淀,排查不直观。
/**
 * 方案二:业务消费端自行控制重试阈值,超限直接归档,不丢去死信队列
 * 原理:消费时读取已重试次数reconsumeTimes,达到上限主动签收成功,停止重试
 * 适用:简单业务,不需要单独死信消费者,代码内直接兜底失败消息
 * @throws Exception
 */
@Test
public void retryConsumer2() throws Exception {
    // 业务正常消费分组
    DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("retry-consumer-group");
    consumer.setNamesrvAddr(MqConstant.NAME_SRV_ADDR);
    // 订阅业务正常Topic
    consumer.subscribe("retryTopic", "*");
    consumer.registerMessageListener(new MessageListenerConcurrently() {
        @Override
        public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
            MessageExt messageExt = msgs.get(0);
            System.out.println("消息接收时间:" + new Date());
            try {
                // 执行业务数据库逻辑(模拟报错:除零异常)
                handleDb();
            } catch (Exception e) {
                // 获取当前消息已经重试多少次
                int reconsumeTimes = messageExt.getReconsumeTimes();
                // 阈值:重试≥3次不再重试
                if (reconsumeTimes >= 3) {
                    // 超限:落地日志、入库、发告警人工处理
                    System.out.println("重试次数过大,归档记录,停止重试");
                    // 返回成功,Broker清除消息,不会进入死信队列
                    return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
                }
                // 不足3次,告知Broker稍后重试
                return ConsumeConcurrentlyStatus.RECONSUME_LATER;
            }
            // 业务无异常,正常签收
            return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
        }
    });
    consumer.start();
    System.in.read();
}

/** 模拟业务数据库异常 */
private void handleDb() {
    // 10/0 抛出算术异常,触发消费失败重试逻辑
    int i = 10 / 0;
}

10.2 RocketMQ 死信消息

当消费重试到达阈值以后,消息不会被投递给消费者了,而是进入了死信队列

  • 当一条消息初次消费失败, RocketMQ 会自动进行消息重试,达到最大重试次数后,若消费依然失败,则表明消费者在正常情况下无法正确地消费该消息。此时,该消息不会立刻被丢弃,而是将其发送到该消费者对应的特殊队列中,这类消息称为死信消息(Dead-Letter Message),存储死信消息的特殊队列称为死信队列(Dead-Letter Queue),死信队列是死信Topic下分区数唯一的单独队列。
  • 如果产生了死信消息,对应的ConsumerGroup的死信Topic名称为%DLQ%ConsumerGroupName,死信队列的消息将不会再被消费。
    • 可以利用 RocketMQ Admin 工具或者 RocketMQ Dashboard 上查询到对应死信消息的信息。

    • 也可以监听死信队列,进行自己的业务上的逻辑,写日志、通知人工处理

消息生产者

@Test
public void testDeadMsgProducer() throws Exception {
    // 生产者分组
    DefaultMQProducer producer = new DefaultMQProducer("dead-group");
    producer.setNamesrvAddr("localhost:9876");
    producer.start();
    // 创建消息:Topic=dead-topic,消息体字符串字节数组
    Message message = new Message("dead-topic", "我是一个死信消息".getBytes());
    // 同步发送消息
    producer.send(message);
    // 关闭生产者释放资源
    producer.shutdown();
}

消息消费者设置全局最大重试次数,失败自动转死信

@Test
@Test
public void testDeadMsgConsumer() throws Exception {
    DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("dead-group");
    consumer.setNamesrvAddr("localhost:9876");
    consumer.subscribe("dead-topic", "*");
    // 全局配置:单条消息最大重试2次
    // 重试2次后依旧失败,RocketMQ自动将消息转入死信队列 %DLQ%dead-group
    consumer.setMaxReconsumeTimes(2);
    consumer.registerMessageListener(new MessageListenerConcurrently() {
        @Override
        public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
            System.out.println("收到业务消息:" + msgs);
            // 固定返回重试,每次都会重新投递,耗尽次数自动进DLQ
            return ConsumeConcurrentlyStatus.RECONSUME_LATER;
        }
    });
    consumer.start();
    System.in.read();
}

死信消费者

注意权限问题

/**
 * 监听死信队列专用消费者
 * 底层规则:消息重试超过maxReconsumeTimes阈值,自动移入 %DLQ%消费组名 死信Topic
 * 普通业务消费者不会订阅DLQ,必须单独写消费者监听处理积压失败消息
 */
@Test
public void testDeadMq() throws  Exception{
    DefaultMQPushConsumer consumer = new DefaultMQPushConsumer("dead-group");
    consumer.setNamesrvAddr("localhost:9876");
    // 订阅死信Topic,格式固定:%DLQ% + 原业务消费者分组名称
    consumer.subscribe("%DLQ%dead-group", "*");
    consumer.registerMessageListener(new MessageListenerConcurrently() {
        @Override
        public ConsumeConcurrentlyStatus consumeMessage(List<MessageExt> msgs, ConsumeConcurrentlyContext context) {
            System.out.println("捕获死信消息:" + msgs);
            // 人工处理、归档、修复后签收,消息彻底清理
            return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
        }
    });
    consumer.start();
    // 阻塞等待消息
    System.in.read();
}

控制台显示

十一、顺序消息:保证同一业务键局部有序

RocketMQ 的顺序消息不是默认全局有序,而是基于 MessageQueue 或 MessageGroup 的局部顺序。 例如同一个订单的状态必须按顺序处理:

订单创建
   ↓
订单支付
   ↓
订单发货
   ↓
订单完成

如果四条消息进入不同 Queue,并由不同消费者并行处理,就可能先收到“订单发货”,再收到“订单支付”。

正确做法是使用 orderId 作为顺序键,让同一订单的消息进入同一 Queue:

/**
 * syncSendOrderly:同步有序发送API
 * 语法:syncSendOrderly(topicTag, payload, hashKey)
 * 原理:依靠第三个参数hashKey做哈希取模路由,相同hashKey消息固定投递到同一个消息队列
 * 作用:保证同一订单的所有状态变更消息进入同一队列,消费者顺序消费,避免状态错乱
 */
rocketMQTemplate.syncSendOrderly(
        // 目的地格式 Topic:Tag,Tag过滤只接收订单状态变更消息
        "order-event-topic:ORDER_STATUS_CHANGED",
        // 消息体:订单变更事件实体,存放完整订单业务数据
        event,
        // 路由哈希键:orderId
        // 相同orderId哈希值一致,会分发到同一个Queue,实现局部有序
        event.getOrderId().toString()
);

顺序键相同的消息会稳定路由到同一 Queue,然后消费者使用顺序消费模式:

import org.apache.rocketmq.spring.annotation.ConsumeMode;
import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;
import org.apache.rocketmq.spring.core.RocketMQListener;
import org.springframework.stereotype.Component;

@Component
// RocketMQ消费者注解,声明监听配置
@RocketMQMessageListener(
        topic = "order-event-topic",                // 监听的主主题
        selectorExpression = "ORDER_STATUS_CHANGED",// Tag过滤表达式,只消费该标签消息
        consumerGroup = "order-state-consumer",     // 消费者分组,同组负载均衡,不同组各自消费全量消息
        consumeMode = ConsumeMode.ORDERLY           // 消费模式:ORDERLY 顺序消费;CONCURRENT 并发消费
)
/**
 * 订单状态变更有序消费者
 * 原理:生产者用orderId路由到同一队列,消费者单线程串行消费单个队列消息
 * 作用:严格按消息发送顺序更新订单状态,防止出现「先收到已完成、后收到待支付」错乱
 */
public class OrderStateConsumer
        implements RocketMQListener<OrderStatusChangedEvent> {

    // 订单业务服务,处理状态更新逻辑
    private final OrderService orderService;

    // 构造器注入
    public OrderStateConsumer(OrderService orderService) {
        this.orderService = orderService;
    }

    /**
     * 消息消费入口方法
     * @param event 自动反序列化后的订单变更事件
     * 有序消费约束:上一条消息处理完成、方法执行完毕,才会拉取下一条同队列消息
     */
    @Override
    public void onMessage(OrderStatusChangedEvent event) {
        // 执行订单状态变更业务逻辑
        orderService.changeStatus(
                event.getOrderId(),
                event.getTargetStatus()
        );
    }
}

RocketMQ 官方说明,顺序消息需要保证生产顺序、存储顺序和消费顺序;同一个消息组内保证顺序,不同消息组之间不保证相互顺序。

顺序消息有几个重要边界:

同一个 Queue 内有序,不等于整个 Topic 全局有序
同一业务键必须使用相同顺序键
生产者并发发送同一业务键时,业务侧也要保证发送先后关系
顺序消费失败会阻塞当前 Queue 后续消息
Queue 越少,全局顺序越强,但并发能力越低

如果某条顺序消息永久失败,会阻塞后续同 Queue 消息,因此业务状态机、异常分类和人工补偿非常重要。


十二、延时消息与定时任务的区别

延时消息表示消息发送后不会立即被消费者看到,而是在指定延迟时间到达后才可消费。官方消息模型将 Delay 定义为:消息在延迟时间到期后才对消费者可见。

典型场景:

订单创建后 30 分钟未支付自动取消
文件删除后 7 天彻底清理
验证码发送 1 分钟后检查状态
任务失败后延迟 10 分钟重试
会议开始前发送提醒

传统 Spring Starter 中常见按延时级别发送:

rocketMQTemplate.syncSend(
        "order-delay-topic:CHECK_PAYMENT",
        message,
        3000,
        delayLevel
);

不同 RocketMQ 版本和客户端对延时级别、定时时间的支持方式不同,必须按当前服务端和客户端文档配置,不能把旧版固定级别代码直接当成所有版本通用写法。

延时消息与定时任务的区别:

定时任务:
    按固定时间或周期扫描
    例如每天凌晨清理数据

延时消息:
    每个业务事件产生自己的未来触发时间
    例如每个订单创建后 30 分钟检查一次

订单超时取消使用延时消息时仍要检查当前状态:

@Transactional
public void cancelIfUnpaid(Long orderId) {
    Order order = orderMapper.selectById(orderId);

    if (order == null) {
        return;
    }

    if (!OrderStatus.CREATED.equals(order.getStatus())) {
        // 已支付、已取消或已完成,不再处理
        return;
    }

    orderMapper.cancelOrder(orderId);
}

因为延时消息可能重复,订单也可能在延迟期间已经支付。


十三、事务消息:解决本地事务与消息发送的一致性

事务消息用于解决:

数据库事务成功,但消息没发送
消息发送成功,但数据库事务失败

RocketMQ 事务消息不是让 RocketMQ 直接参与数据库两阶段提交,也不是分布式数据库事务。它通过 Half Message、本地事务执行结果和事务状态回查实现最终一致性。

官方事务消息流程可以概括为:

1. Producer 发送 Half Message
2. Broker 保存消息,但暂不投递给 Consumer
3. Broker 返回 Half Message 成功
4. Producer 执行本地事务
5. Producer 向 Broker 提交 Commit 或 Rollback
6. Commit 后消息才对 Consumer 可见
7. 状态未知时,Broker 回查 Producer

官方文档明确描述:Broker 先保存一条不可投递的 Half Message,Producer 再执行本地事务,并向 Broker提交 Commit 或 Rollback;如果状态不明确,Broker 可以进行事务状态检查。

Spring 事务消息发送示意:


/**
 * 订单事务消息生产者
 * 作用:发送事务半消息,绑定本地创建订单DB事务,保证「建订单库写入」和「MQ消息投递」原子一致
 * 原理:2PC事务消息,先发半消息,再执行本地事务,根据本地事务结果决定消息对外可见/丢弃
 */
@Service
@RequiredArgsConstructor // Lombok:final字段构造器注入rocketMQTemplate
public class OrderTransactionProducer {

    // RocketMQ操作模板,封装事务消息发送能力
    private final RocketMQTemplate rocketMQTemplate;

    /**
     * 创建订单入口,发送事务消息
     * @param command 创建订单请求载体,携带订单所有业务参数
     */
    public void createOrder(CreateOrderCommand command) {
        // 1、构建标准Spring Message消息
        Message<CreateOrderCommand> message =
                MessageBuilder.withPayload(command) // Payload:消息主体,消费者接收解析
                        // MQ内置KEY索引字段:requestId作为全局唯一标识
                        // 用途1:写入IndexFile哈希索引,控制台可按requestId检索消息
                        // 用途2:事务回查时唯一查询条件,精准匹配本地订单记录
                        .setHeader(
                                RocketMQHeaders.KEYS,
                                command.getRequestId()
                        )
                        .build();

        // 2、发送事务消息核心方法
        // 参数1:目的地 Topic:Tag 格式,ORDER_CREATED标记「订单创建」事件
        // 参数2:完整消息对象
        // 参数3:自定义透传参数,会原样传入 executeLocalTransaction 的argument参数
        // 返回值:TransactionSendResult 封装本次本地事务执行结果
        TransactionSendResult result =
                rocketMQTemplate.sendMessageInTransaction(
                        "order-transaction-topic:ORDER_CREATED",
                        message,
                        command
                );

        // 3、判断本地事务执行状态,事务回滚则主动抛出业务异常
        if (result.getLocalTransactionState()
                == LocalTransactionState.ROLLBACK_MESSAGE) {
            throw new IllegalStateException(
                    "订单事务消息执行失败"
            );
        }
    }
}

本地事务监听器:

import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.apache.rocketmq.spring.annotation.RocketMQTransactionListener;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionListener;
import org.apache.rocketmq.spring.core.RocketMQLocalTransactionState;
import org.apache.rocketmq.spring.support.RocketMQHeaders;
import org.springframework.messaging.Message;

/**
 * RocketMQ事务消息监听器
 * 两个核心方法:
 * 1. executeLocalTransaction:半消息发送成功后同步执行本地DB事务
 * 2. checkLocalTransaction:Broker定时回查,生产者宕机/网络丢失时兜底校验事务状态
 */
@RocketMQTransactionListener // 标识事务监听,自动绑定对应事务生产者组
@RequiredArgsConstructor
@Slf4j // 自动生成log日志对象
public class OrderTransactionListener
        implements RocketMQLocalTransactionListener {

    // 订单业务服务:封装创建订单数据库逻辑
    private final OrderService orderService;
    // 订单Mapper:回查时查询数据库订单记录
    private final OrderMapper orderMapper;

    /**
     * 阶段1:半消息持久化Broker后,同步回调执行本地事务
     * @param message MQ原始消息对象
     * @param argument sendMessageInTransaction第三个参数,透传CreateOrderCommand
     * @return 事务状态:COMMIT提交消息 / ROLLBACK丢弃消息 / UNKNOWN状态未知等待回查
     */
    @Override
    public RocketMQLocalTransactionState executeLocalTransaction(
            Message message,
            Object argument
    ) {
        // 取出发送时透传的订单创建参数
        CreateOrderCommand command =
                (CreateOrderCommand) argument;

        try {
            // 执行本地数据库事务:生成订单、扣库存、落订单表等
            orderService.createOrder(command);
            // 本地事务成功,通知Broker把半消息转正,投递到业务Topic供消费者消费
            return RocketMQLocalTransactionState.COMMIT;
        } catch (Exception e) {
            log.error(
                    "订单本地事务执行失败, requestId={}",
                    command.getRequestId(),
                    e
            );
            // DB事务异常,通知Broker删除半消息,下游永远收不到该订单消息
            return RocketMQLocalTransactionState.ROLLBACK;
        }
    }

    /**
     * 阶段2:事务补偿回查接口
     * 触发条件:生产者发送半消息后宕机/网络断连,Broker长期未收到commit/rollback指令
     * 逻辑:通过消息头里唯一requestId查询数据库,判断订单事务是否实际创建成功
     */
    @Override
    public RocketMQLocalTransactionState checkLocalTransaction(
            Message message
    ) {
        // 从内置KEY头取出全局唯一请求ID
        String requestId = String.valueOf(
                message.getHeaders().get(
                        RocketMQHeaders.KEYS
                )
        );

        // 根据requestId查询订单表,判断本地事务是否执行成功
        Order order = orderMapper.selectByRequestId(
                requestId
        );

        if (order != null) {
            // 数据库存在订单记录 → 本地事务已完成,通知Broker提交消息
            return RocketMQLocalTransactionState.COMMIT;
        }
        // 无订单记录 → 本地事务回滚/未执行,Broker删除半消息
        return RocketMQLocalTransactionState.ROLLBACK;
    }
}

事务回查不能依赖内存变量,因为 Producer 可能重启或请求可能被路由到另一实例。必须根据业务唯一键查询持久化状态:

根据 requestId 查询订单是否存在
根据 paymentNo 查询支付记录是否成功
根据 eventId 查询本地事务日志

事务消息适合:

创建订单后发送订单事件
支付成功后发送支付结果
账户变更后发送记账事件
文件状态变化后触发异步处理

事务消息仍然不能消除重复消费,因此 Consumer 依旧必须幂等。


十四、消息积压、消费能力与流量控制

消息积压指 Broker 中尚未被某消费组处理的消息越来越多:

生产速度:每秒 10,000 条
消费速度:每秒 6,000 条
净积压:每秒 4,000 条

只要生产速度长期大于消费速度,积压一定持续增长。

14.1 常见积压原因

消费者实例太少
Topic Queue 数量不足
单条消息处理耗时过长
数据库慢 SQL
第三方接口超时
消费者频繁抛异常重试
顺序消息被某条异常消息阻塞
消息体过大
消费线程池配置不合理
消费者发生 Full GC
Broker、网络或磁盘性能异常

14.2 处理积压的正确顺序

首先确认是生产突增还是消费下降:

生产 TPS
消费 TPS
消费延迟
最大积压量
失败重试数
消费者实例数
数据库响应时间
Broker 磁盘利用率

然后定位消费者单条处理耗时。如果业务处理是 I/O 密集型,可以适当增加消费线程或实例;如果瓶颈在数据库,盲目增加消费者只会把数据库压垮。

扩容消费者前要确认 Queue 数量:

Topic 只有 4 个 Queue
Consumer Group 启动 20 个实例

大部分实例不会获得有效 Queue,水平扩容收益很小。要提升并行度,通常需要同时规划 Queue 数量和消费者实例数量。

14.3 消费者批处理

对于写数据库、写 Elasticsearch 或调用批量接口的场景,可以考虑批量消费:

逐条写数据库:
1000 条消息 → 1000 次数据库请求

批量写数据库:
1000 条消息 → 10 次批量请求

但批量处理要明确部分成功时怎么办。不能因为 100 条中 1 条失败,就无脑重复执行另外 99 条非幂等业务。

14.4 背压与限流

消费者不能无限制从 Broker 拉取并堆在 JVM 内存中。客户端通常通过拉取批量、消费线程池、本地缓存阈值等参数控制流量。

企业配置原则是:

根据下游承载能力设置消费并发
根据消息大小控制本地缓存
根据平均耗时评估线程数
对第三方调用增加超时、限流和熔断
对数据库写入使用批量和索引优化

十五、集群部署、高可用与生产环境设计

本地开发可以使用:

1 个 NameServer
1 个 Broker

但生产环境不能把这种单节点结构直接上线。

典型生产结构:

3 个 NameServer

Broker Group A
    Broker A 主节点
    Broker A 副本

Broker Group B
    Broker B 主节点
    Broker B 副本

多个 Producer 实例
多个 Consumer 实例
Dashboard
Prometheus + Grafana
集中日志系统

15.1 多 Broker 的作用

多个 Broker 可以:

分散 Topic Queue
提高集群吞吐量
避免单节点成为容量瓶颈
降低单机故障影响
实现水平扩展

15.2 主从复制

Broker 主节点接收写入,从节点保存副本。复制可以分为同步和异步。

同步复制:

Master 写入
    ↓
等待 Slave 确认
    ↓
向 Producer 返回成功

可靠性更高,但发送延迟更大。

异步复制:

Master 写入
    ↓
立即返回成功
    ↓
后台复制给 Slave

性能更高,但 Master 在复制完成前发生不可恢复故障,可能损失少量消息。

15.3 Broker 扩容

新增 Broker 后,需要让新 Broker 向相同 NameServer 注册,并调整 Topic 的 Queue 分布。官方 Dashboard 文档也提供扩展 Broker 和调整 Topic Queue 的运维入口。

扩容不是只启动一台机器,还要评估:

Topic Queue 是否迁移或新增
Producer 路由是否刷新
Consumer 是否发生 Rebalance
顺序键映射是否变化
磁盘容量与水位
网络带宽
消息副本
监控和告警

15.4 消息保留与清理

RocketMQ 的消息通常按存储时间和磁盘策略清理,不是“消费者消费后立刻删除”。不同 Consumer Group 可以按各自进度读取同一条消息,因此消息不能因某一个消费组成功就立即物理删除。

官方文档将消息存储和清理作为 Broker 的独立策略管理。

这意味着 RocketMQ 与普通任务队列的直觉不同:

消费成功:
    更新消费组 Offset

物理删除:
    后台按文件保留和磁盘策略清理 CommitLog

十六、监控、日志与排障

RocketMQ 上线后至少要监控:

Broker 是否存活
NameServer 是否存活
Producer 发送成功率
Producer 发送耗时
Consumer 消费 TPS
Consumer 失败率
Consumer 重试次数
消息积压量
消费延迟
死信消息数量
Broker 磁盘使用率
CommitLog 写入耗时
主从复制延迟
JVM 堆、GC 和线程数
网络连接数

RocketMQ 5.x 官方提供 Prometheus 格式指标,覆盖 Broker、Producer 和 Consumer 指标;相关指标能力在 5.1.0 后逐步引入。

16.1 消息排障四个关键标识

日志中至少记录:

eventId
businessId
messageId
traceId

Producer 日志:

log.info(
        "消息发送成功, eventId={}, businessId={}, msgId={}, topic={}",
        eventId,
        businessId,
        sendResult.getMsgId(),
        topic
);

Consumer 日志:

log.info(
        "消息消费开始, eventId={}, businessId={}, consumerGroup={}",
        eventId,
        businessId,
        consumerGroup
);

出现问题时按以下顺序排查:

1. 业务数据库是否提交成功
2. Producer 是否生成事件
3. Producer 是否发送成功
4. Broker 是否能按 Key 查到消息
5. 消费组是否订阅正确 Topic 和 Tag
6. 消费组 Offset 是否推进
7. 消息是否进入重试队列
8. 消息是否进入死信队列
9. Consumer 是否异常或频繁重启
10. 下游数据库、Redis、ES 是否异常

16.2 常见错误

Topic 或 Tag 不一致

Producer:

file-event-topic:FILE_UPLOADED

Consumer:

topic = files-event-topic
selectorExpression = FILE_UPLOAD

字符串不同,消费者自然收不到。

Consumer Group 使用错误

两个完全不同的业务消费者误用同一个 Group,会互相竞争消息。

Consumer 未启动或 Bean 未被扫描

监听器类没有 @Component,或者不在 Spring Boot 启动类扫描范围内。

NameServer 与 Broker 地址问题

客户端能访问 NameServer,但 Broker 注册的是容器内部 IP 或错误网卡地址,导致客户端获取路由后无法连接 Broker。Docker、云服务器、多网卡和 VPN 环境尤其常见。

消费异常被吞掉

业务失败但监听方法正常返回,RocketMQ 会认为消费成功。

消息反序列化失败

Producer 和 Consumer 使用不同 DTO 结构、包名、字段类型或序列化方式。事件 DTO 要有版本管理,并尽量保持向后兼容。


十七、企业级消息设计原则

17.1 消息是事件,不是远程方法调用参数

不推荐:

{
  "method": "updateFile",
  "args": [...]
}

推荐:

{
  "eventId": "evt-001",
  "eventType": "FILE_UPLOADED",
  "fileId": 1001,
  "userId": 2001,
  "occurredAt": "2026-07-11T20:30:00",
  "schemaVersion": 1
}

事件描述“发生了什么”,消费者自己决定如何响应。

17.2 消息体不能依赖数据库瞬时状态

只发送:

{
  "orderId": 10001
}

然后消费者查数据库,结构简单,但可能遇到数据库记录已变化或被删除。

发送完整快照:

{
  "orderId": 10001,
  "userId": 20001,
  "amount": 99.00,
  "status": "PAID"
}

可以减少查询,但消息更大,也可能包含敏感数据。

企业中通常根据一致性和数据量折中:

发送业务主键
+ 消费者真正需要的稳定字段
+ 事件发生时状态

17.3 Schema 要兼容升级

Producer 增加字段时,旧 Consumer 应能忽略;删除或修改字段类型时,要考虑旧消息仍然可能存在于 Broker 中。

推荐:

新增可选字段
保持旧字段语义
使用 schemaVersion
灰度升级 Consumer,再升级 Producer
避免直接修改字段类型

17.4 Topic 和 Consumer Group 命名规范

例如:

Topic:
networkdisk-file-event
networkdisk-user-event
networkdisk-share-event

Consumer Group:
networkdisk-ai-file-index-consumer
networkdisk-search-file-sync-consumer
networkdisk-notification-file-consumer

命名中应体现业务域和消费用途,不要使用:

test-topic
topic1
consumer-group
default-group

17.5 大消息不要直接放 MQ

图片、视频、文件二进制不应该直接作为普通消息体发送。应该把文件保存到对象存储或文件系统,然后消息中只传:

fileId
objectKey
bucket
metadata

大消息会增加 Broker 网络、磁盘、内存和消费重试成本。


十八、RocketMQ 与 Kafka、RabbitMQ 的区别

RocketMQ

特点:

面向业务消息和事件驱动
支持事务消息
支持顺序消息
支持延时消息
支持消费重试和死信
Topic、Queue、ConsumerGroup 模型清晰
在 Java 和微服务业务系统中使用广泛

适合订单、支付、库存、通知、异步任务和最终一致性。

Kafka

更偏向高吞吐事件流、日志流和数据管道:

日志采集
CDC
大数据流处理
埋点
实时计算
长期事件流

Kafka 以 Partition 和 Consumer Group 为核心,生态在流处理和大数据领域更成熟。

RabbitMQ

基于 AMQP,Exchange、Queue、RoutingKey 路由模型灵活:

复杂路由
中小规模业务消息
工作队列
灵活确认机制

选择时不能只比较“谁性能高”,而应比较:

消息模型
事务和延时需求
顺序要求
吞吐量
运维经验
团队技术栈
云服务支持
现有基础设施

十九、RocketMQ 高频面试题

1. RocketMQ 有什么作用?

主要解决异步、解耦和削峰问题,同时支持最终一致性、延时任务、顺序处理、事务消息和事件驱动。它适合不要求同步返回结果的跨服务通信,但不能完全替代 Dubbo、HTTP 等同步调用。

2. NameServer 和 Broker 有什么区别?

NameServer 保存 Broker 和 Topic 路由,不存业务消息;Broker 接收、持久化和投递消息,并维护消费进度、重试、死信和索引。Producer 和 Consumer 从 NameServer 获取路由后直接连接 Broker。

3. Topic 和 MessageQueue 有什么区别?

Topic 是消息的逻辑业务分类,MessageQueue 是 Topic 下的物理并行和顺序单元。多个 Queue 提高吞吐和消费并行度,但 RocketMQ 通常只保证单 Queue 内有序。

4. Consumer Group 有什么作用?

同一个 Group 的多个实例共同分担消息,实现负载均衡和水平扩展;不同 Group 之间相互独立,每个 Group 都可以消费同一 Topic 的完整消息。

5. RocketMQ 为什么快?

主要依靠 CommitLog 顺序追加写、ConsumeQueue 轻量索引、内存映射、Page Cache、批量刷盘以及多 Queue 并行,而不是简单地“不写磁盘”。

6. CommitLog 和 ConsumeQueue 的区别?

CommitLog 保存完整消息,是物理存储;ConsumeQueue 保存 Topic 某个 Queue 对应的逻辑索引,通过物理偏移定位 CommitLog 中的消息。

7. RocketMQ 如何保证消息不丢?

要分三个阶段:

生产阶段:
同步发送、失败重试、本地消息表、事务消息、补偿任务

Broker 阶段:
持久化、同步刷盘、主从复制、高可用部署

消费阶段:
业务成功后再确认、失败重试、死信告警、消费幂等

任何单一配置都不能独立保证端到端绝对可靠。

8. 为什么会重复消费?

发送响应超时、消费成功但 ACK 丢失、消费者宕机、Rebalance、重试和人工重置 Offset 都可能导致重复。消费者应使用事件 ID、数据库唯一键、状态机和本地事务实现幂等。

9. 如何保证顺序消费?

Producer 根据业务键选择固定 Queue,Consumer 使用顺序消费模式。顺序只在同一 Queue 或消息组内成立,不同 Queue 之间不保证全局顺序。

10. 事务消息原理是什么?

Producer 先发送 Half Message,Broker 保存但不投递;Producer 执行本地事务后返回 Commit 或 Rollback;状态未知时 Broker 回查生产者的本地事务状态。它实现的是消息与本地事务的最终一致性,不是数据库 XA 两阶段提交。

11. 消息积压怎么处理?

先判断生产 TPS 是否突增、消费 TPS 是否下降,再检查消费者异常、慢 SQL、下游超时、线程池、Queue 数量和实例数量。扩容消费者前必须确认 Queue 是否足够,不能盲目增加实例。

12. 重试与死信有什么区别?

重试是消费失败后由 Broker 再次投递;超过最大重试次数仍失败的消息进入死信队列,需要监控、人工分析和补偿,不能继续无限自动重试。

13. 集群消费和广播消费有什么区别?

集群消费中,一条消息在同一 Group 内通常只由一个实例处理;广播消费中,同一 Group 的每个实例都会处理一条消息。普通业务使用集群消费,本地缓存刷新等场景可以使用广播消费。

14. 同步发送成功是否代表消费成功?

不代表。同步发送成功通常只表示 Broker 接收并按当前配置保存消息,消费者什么时候处理、是否处理成功,是后续独立过程。

15. RocketMQ 能保证 Exactly Once 吗?

消息中间件很难单独保证业务效果严格只发生一次。RocketMQ 通常按照至少一次可靠投递设计,业务层通过唯一事件 ID、唯一约束、状态机和事务实现最终幂等。


二十、核心知识主线

RocketMQ 可以用下面一条主线整体理解:

业务服务产生事件
        ↓
Producer 从 NameServer 查询路由
        ↓
Producer 选择 Topic 下的 MessageQueue
        ↓
Broker 将完整消息顺序写入 CommitLog
        ↓
Broker 构建 ConsumeQueue 和 Key 索引
        ↓
Consumer Group 获取并分配 Queue
        ↓
Consumer 拉取消息并执行业务
        ↓
成功后推进消费 Offset
        ↓
失败则进入重试流程
        ↓
超过最大次数进入死信队列

企业级使用 RocketMQ 的重点不是会写 syncSend()@RocketMQMessageListener,而是理解以下几组关系:

异步解耦与数据最终一致性
可靠投递与重复消费
消费重试与业务幂等
MessageQueue 数量与消费并行度
局部顺序与系统吞吐量
同步刷盘与性能
主从复制与可用性
事务消息与本地事务回查
消息积压与下游承载能力

真正可靠的 RocketMQ 系统通常由以下部分共同构成:

规范的 Topic 和 Consumer Group
稳定的消息 Schema
全局唯一事件 ID
Producer 发送确认
本地消息表或事务消息
Broker 多副本和持久化
Consumer 业务幂等
失败重试与死信处理
积压和延迟监控
定时对账与业务补偿
完整的 TraceId、Key 和 MessageId 日志

RocketMQ 提供的是可靠消息基础设施,但端到端业务可靠性仍然需要数据库事务、幂等设计、状态机、监控和补偿机制共同完成。