前言

在分布式系统开发中,消息队列是实现高性能、高可靠应用的关键组件。RocketMQ​ 作为一款诞生于生产实践、并历经“双十一”万亿级流量考验的消息中间件,以其卓越的性能和丰富的功能,成为众多企业级应用的首选。


一、RocketMQ简介

在构建高并发、可扩展的分布式系统时,组件之间的可靠异步通信是至关重要的基石。面对海量消息洪峰,如何确保每一条数据都能被高效、有序、不丢失地处理?Apache RocketMQ 作为一款诞生自阿里业务的顶级开源分布式消息中间件,以其万亿级消息吞吐、低延迟、高可用的卓越能力,给出了完美的答案。
本篇将带你深入浅出地理解RocketMQ的核心架构、设计理念与最佳实践。无论你是希望解决线上系统的解耦、削峰填谷、异步处理难题,还是想要探索其事务消息、顺序消息、延时消息等高级特性,这里都将为你提供一个清晰的路线图。
官网地址:https://rocketmq.apache.org/zh/

二、RocketMQ框架

在这里插入图片描述
这张架构图清晰展示了RocketMQ作为分布式消息系统的核心工作流程:生产者(Producer)和消费者(Consumer)首先从轻量级注册中心(NameServer)获取路由信息,随后生产者将消息发送至Broker集群进行存储(采用主从模式保证高可用),最终由消费者组以订阅模式从Broker拉取并消费消息,实现了从生产、路由、存储到消费的完整、松耦合的异步通信。

三、RocketMQ基本概念

1、NameServer

轻量级注册中心,管理Broker路由信息,实现服务发现(Producer/Consumer通过NameServer查找Broker),无状态设计,节点间不通信,高可用通过多节点实现。

2、Broker

消息存储和转发核心组件,接收生产者消息、存储消息、处理消费者拉取请求,支持主从架构(Master-Slave)实现高可用,负责消息持久化、过滤、查询等。

3、生产者Producer

消息发送方,负责产生消息,支持同步、异步、单向发送,支持事务消息,可集群部署,通过NameServer发现Broker。

4、生产者组ProducerGroup

生产者组,用于事务消息,事务回查会查询生产者组内任意一台。

5、消费者Consumer

消息接收方,负责消费消息,支持集群模式(负载均衡)和广播模式,提供Push和Pull两种消费方式
支持顺序消费、并发消费。

6、消费者组ConsumerGroup

消费者组,相同逻辑的消费者集合
集群模式下,组内消费者分摊消费消息
广播模式下,组内每个消费者都收到全量消息

7、主题Topic

消息的逻辑分类,生产者向指定Topic发送消息,消费者订阅感兴趣的Topic进行消费,一个Topic可被多个消费者组订阅。

8、队列Queque

分为ConsumerQuequeMessageQueque,ConsumerQueque是Topic的物理存储单元,存储消息在Commitlog的索引,和tag哈希信息。MessageQueque是Topic的逻辑区分,只存储队列id、主题和broker信息。一个Topic包含多个MessageQuequeQueue,一个MessageQueque下有多个ConsumerQueque,消息在ConsumerQueque顺序存储,Queue是并行消费和负载均衡的基本单位。区分读写权限,写权限标记为可写,读权限标记对消费者是否可见。默认每个Topic有4个读写队列。

9、消息Message

消息是RocketMQ 中数据传输的最小单位。有以下几种类型:
普通消息:无特殊特性
顺序消息:保证局部顺序(同一队列内FIFO)
事务消息:分布式事务支持,两阶段提交
定时/延时消息:指定时间或延迟后投递
批量消息:一次发送多条消息提高吞吐

10、标签Tag和Key

Tag:二级消息分类,用于细粒度过滤,消费者可通过Tag订阅同一Topic下的部分消息,一条消息只能设置一个tag。
Key:用于消息索引、追踪、标记等,一条消息可以有多个。

四、生产者原理分析

1、基本流程

RocketMQ生产者发送消息的核心流程是:生产者启动,设置生产者组名,检查配置,初始化netty网络环境,从NameServer获取Topic路由信息(包含了broker信息),发送消息时通过负载均衡策略选择一个MessageQueue,将消息通过Netty网络层发送到对应的Broker,Broker处理完成后返回响应结果;如果发送失败,生产者会根据重试机制自动更换Broker重试,最终将发送结果返回给应用程序,整个过程通过心跳机制、故障规避和本地缓存优化来保证高可用性和高性能。
在这里插入图片描述

2、定时任务

生产者会开启定时任务自动感知集群变化、保持连接健康、实现高可用的关键机制,确保生产者在分布式环境中能稳定运行。
每30秒从NameServer拉取路由
每30秒清理下线的Broker
每30秒向所有Broker发送心跳
源码位置:MQClientInstance.startScheduledTask

3、缓存机制

生产者在发送消息之前会尝试从本地缓存获取路由信息,如果本地缓存没有,从远程Nameserver获取并更新本地缓存,能够在高并发、低延迟的场景下保持高性能。
在这里插入图片描述
在这里插入图片描述

4、负载均衡

RocketMQ的生产者负载均衡机制主要涉及消息队列(MessageQueue)的选择策略,以确保消息在多个Broker和队列间均匀分布。以下是其核心机制和工作原理:
核心机制:队列选择策略
在这里插入图片描述

生产者发送消息时,需要通过MessageQueueSelector选择目标队列。RocketMQ提供以下内置策略:(1)轮询策略(默认)
机制:按顺序轮流选择主题下的所有队列。如果开启sendLatencyFaultEnable 延迟故障规避,选不到队列会选择其他队列。
优点:保证消息均匀分布到所有队列,实现负载均衡。
(2)SelectMessageQueueByHash 哈希策略
机制:根据消息的Key(如订单ID)计算哈希值,固定映射到特定队列。
优点:保证相同Key的消息始终进入同一队列,实现顺序消息。
(3)SelectMessageQueueByRandom 随机策略
机制:随机选择一个队列。
适用场景:简单场景,但可能分布不均。
(4)SelectMessageQueueByMachineRoom 就近策略
机制:优先选择与生产者同机房或延迟低的Broker上的队列。
优点:减少网络延迟,提升性能。

5、重试机制

RocketMQ的生产者重试机制主要发生在消息发送过程中,当消息发送失败时,RocketMQ的生产者会自动重试

    # retryTimesWhenSendFailed	同步发送失败的话,rocketmq内部重试多少次	int	2
    # retryTimesWhenSendAsyncFailed	异步发送失败的话,rocketmq内部重试多少次	int	2

6、三种基础发送模式

RocketMQ的发送模式有这三种基础模式

public enum CommunicationMode {
    SYNC,
    ASYNC,
    ONEWAY;
    private CommunicationMode() {
    }
}

SYNC(同步发送)
默认就是同步发送,调用send()方法时,会在当前线程同步阻塞等待Broker响应,内部通过CountDownLatch或Future实现等待,支持失败自动重试(可配置重试次数)。
ASYNC(异步发送)
非阻塞调用,内部通过线程池处理回调,Netty的I/O线程发送请求,业务线程处理回调
ONEWAY(单向发送)
只将消息放入发送队列,不等待任何响应,不保证消息到达Broker,不保证不重复,无等待,无回调,完全"fire and forget"

7、顺序发送

顺序消息,根据队列策略计算选择的队列(用的是key-hash值计算),只有都往一个队列发送才能保证顺序性。

public class SelectMessageQueueByHash implements MessageQueueSelector {
    public SelectMessageQueueByHash() {
    }

    public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
        int value = arg.hashCode() % mqs.size();
        if (value < 0) {
            value = Math.abs(value);
        }

        return (MessageQueue)mqs.get(value);
    }
}

8、批量发送

这个不做过多介绍,用的Spring-RocketMQ的话,就算批量发送,还是会单条消费,不过倒是有个好处,可以减少网络IO频率。

9、事务消息发送

事务消息是 Apache RocketMQ 提供的一种高级消息类型,支持在分布式场景下保障消息生产和本地事务的最终一致性。
主要流程如下:
在这里插入图片描述
1、生产者将消息发送至Apache RocketMQ服务端。

2、Apache RocketMQ服务端将消息持久化成功之后,向生产者返回Ack确认消息已经发送成功,此时消息被标记为"暂不能投递",这种状态下的消息即为半事务消息。

3、生产者开始执行本地事务逻辑。

4、生产者根据本地事务执行结果向服务端提交二次确认结果(Commit或是Rollback),服务端收到确认结果后处理逻辑如下:

二次确认结果为Commit:服务端将半事务消息标记为可投递,并投递给消费者。

二次确认结果为Rollback:服务端将回滚事务,不会将半事务消息投递给消费者。

5、在断网或者是生产者应用重启的特殊情况下,若服务端未收到发送者提交的二次确认结果,或服务端收到的二次确认结果为Unknown未知状态,经过固定时间后,服务端将对消息生产者即生产者集群中任一生产者实例发起消息回查。 说明 服务端回查的间隔时间和最大回查次数,请参见参数限制。

6、生产者收到消息回查后,需要检查对应消息的本地事务执行的最终结果。

7、生产者根据检查到的本地事务的最终状态再次提交二次确认,服务端仍按照步骤4对半事务消息进行处理。

10、延时发送

在分布式定时调度触发、任务超时处理等场景,需要实现精准、可靠的定时事件触发。使用 Apache RocketMQ 的定时消息可以简化定时调度任务的开发逻辑,实现高性能、可扩展、高可靠的定时触发能力。
典型场景:分布式定时调度
在这里插入图片描述
以电商交易场景为例,订单下单后暂未支付,此时不可以直接关闭订单,而是需要等待一段时间后才能关闭订单。使用 Apache RocketMQ 定时消息可以实现超时任务的检查触发。

Apache RocketMQ 一共支持18个等级的延迟投递(不过新版本5.x已经支持精准投递,能按实际时间去投),具体时间如下:
在这里插入图片描述

五、Broker原理分析

1、存储结构

在这里插入图片描述
先来看看实际存储的基本概念:
1.1、CommitLog(物理存储)
所有主题的消息都按顺序、定长地写入同一个物理文件。这是真正的消息主体存储文件。这种顺序写盘机制能极大提升磁盘IO效率。文件名是20位的数字起始偏移量(如图中的00000000001073741824),表示这个文件中第一条消息的全局物理偏移量。单个文件大小固定(默认为1GB),写满后创建新文件。

1.2、ConsumeQueue(存储队列)
consumequeue/TopicTest/0/表示 Topic为“TopicTest”,MessageQueque队列ID为0​ 的存储队列目录。
这是消息的“索引目录”。它为每个Topic下的每个Message Queue建立一个存储队列文件,是消费者消费进度的直接依据,相当于数据库的id主键索引。

目录下的文件(如00000000000000000000)是索引文件。每个条目固定20字节,包含:
消息在CommitLog中的物理偏移量(8字节)
消息长度(4字节)
消息Tag的哈希码(8字节)

工作流程:
消息写入CommitLog后,会异步生成索引条目,追加到对应的ConsumeQueue文件中。
消费者拉取消息时,先根据消费进度(offset)从ConsumeQueue中找到索引条目。
再根据索引条目中的物理偏移量和长度,到CommitLog中精准读取出完整的消息内容。
这种设计的优势:写入是单一的、顺序的磁盘写入(CommitLog),速度极快。读取时,虽然大部分是随机读,但ConsumeQueue文件小且顺序读取,效率依然很高,并且利用操作系统的PageCache可大幅提升性能

1.3、IndexFile(消息索引)
位于index目录下,文件名是时间戳。
用于支持按Key查询和按时间范围查询消息的命令,构建了消息Key到物理偏移量的哈希索引。可以理解为数据库的二级索引

1.4、配置文件(config目录)
consumerOffset.json:记录各个消费者组的消费进度。
topics.json:记录Broker上创建的主题和队列配置信息。
subscriptionGroup.json:订阅组配置。
(.bak文件是自动备份)

2、页缓存

页缓存是RocketMQ实现高性能、高吞吐的核心机制之一。
页缓存(Page Cache)​ 是操作系统内核实现的一种磁盘数据缓存机制:
当应用程序读取磁盘文件时,内核会将磁盘数据缓存在内存中
当应用程序写入文件时,数据先写入内存缓存,再由内核异步刷回磁盘
对应用程序来说,读写操作似乎都是"内存操作",性能极高

3、刷盘机制

页缓存是将数据写入内存缓存中,最终持久化需要刷到磁盘上,这里就有两种刷盘机制:

同步刷盘(SYNC_FLUSH)
消息写入PageCache后,立即调用fsync()强制刷盘,确认落盘后才返回成功,这样能保证消息不丢失,但是性能比较差一点。

异步刷盘(ASYNC_FLUSH)
消息写入PageCache后立即返回成功,由后台线程定期异步刷盘,性能高,但是可能丢数据

高可靠场景:

# 主节点:同步刷盘 + 同步复制
brokerRole=SYNC_MASTER
flushDiskType=SYNC_FLUSH
syncFlushTimeout=10000  # 10秒超时

# 从节点:可异步刷盘
brokerRole=SLAVE
flushDiskType=ASYNC_FLUSH

性能优先(大数据场景)

# 异步刷盘 + 异步复制
brokerRole=ASYNC_MASTER
flushDiskType=ASYNC_FLUSH
# 降低刷盘频率提升性能
flushIntervalCommitLog=1000  # 1秒刷一次
flushLeastPages=16  # 64KB才刷盘
transientStorePoolEnable=true

还有其他组合,比如同步刷盘+异步复制,异步刷盘+同步复制斟酌使用,反正同步刷盘和同步复制是消息可靠性的保障,异步则是性能吞吐量的能力。

4、零拷贝

零拷贝是RocketMQ实现高性能网络传输的关键技术。它的核心原理是让数据在内核空间中直接完成传输,绕过在应用层内存的多次拷贝。具体来说,当Broker向消费者发送消息时,会利用操作系统的sendfile或FileChannel.transferTo()调用,将CommitLog中消息数据直接从页缓存传输到网卡缓冲区,而无需先将数据拷贝到JVM堆内存,再由应用程序写入Socket。这项技术消除了不必要的CPU拷贝和上下文切换,极大地提升了吞吐量并降低了延迟,是支撑RocketMQ高并发能力的重要基石。说白了就是底层直接传输,不走jvm,减少用户态和内核态CPU拷贝工具。

5、主从复制

RocketMQ的主从复制是高可用架构的核心,通过将Master节点的数据复制到Slave节点,实现数据备份和故障转移。

异步复制(ASYNC_MASTER)
主节点写入成功后立即返回,不等待从节点确认。

# Master配置
brokerRole=ASYNC_MASTER
# Slave配置  
brokerRole=SLAVE

同步复制(SYNC_MASTER)
必须等待至少一个Slave确认后才返回生产者,主从数据强一致。

# Master配置
brokerRole=SYNC_MASTER
# Slave配置
brokerRole=SLAVE

六、消费者原理分析

1.基本流程

在这里插入图片描述
这里流程就介绍push流程,因为默认也是用这个。RocketMQ 的消费流程始于消费者启动后从 NameServer 获取路由信息并连接 Broker,随后通过长轮询机制持续拉取消息。消息根据消费模式(集群模式下由同组消费者分担队列,广播模式下各自消费全量)被分配给对应的消费者实例,并由注册的消息监听器(分为无序的并发处理与严格有序的队列处理)执行业务逻辑。消费成功后,消费进度(Offset)​ 会被定时提交回 Broker 进行持久化管理;若消费失败,消息会进入重试队列,并按照延迟策略进行多次重试,最终失败则转入死信队列。整个流程通过心跳维持连接,并通过流量控制与偏移量管理,确保了消息的可靠投递与弹性处理能力。

2.消费者分类

在这里插入图片描述
一般用PushConsumer,PullConsumer基本不用,官网连示例都没有。
Push模式分析:

核心概念

长轮询机制 消费者启动后立给监听的队列创建PullRequest,往请求队列添加请求,然后while循环不断地take请求队列,向 Broker 发送拉取请求,若请求不满足某些条件(例如当前当前缓存消息超过缓存数量、超过缓存大小,大于消费跨度等),延迟放回请求队列,若broke队列无消息,Broker 挂起请求(默认 15s),开启长轮询,期间有新消息到达立即返回,否则超时后返回空。

ProcessQueue ,ProcessQueue是 RocketMQ 消费端的核心数据结构,它是每个消息队列(MessageQueue)在消费者本地的内存镜像和状态管理器。包含:

2.1、按offset缓存消息
从broker拉取下来的消息会缓存到ProcessQueue 的msgTreeMap,这是一个树结构map,可能会有疑问,为什么树结构?保证消息有序,为什么需要缓存?提升吞吐量,缓存数据利于数据管理,比如消息加锁顺序消费,平衡拉取消费不一致,保护消费者。

2.2、消费进度跟踪
简单说下,进度管理那块再讲。说白了,就是通过 msgTreeMap的最小key和最大可以实时跟踪已拉取但未消费的消息。消费成功后从缓存移除,并异步提交进度。内存 → 本地 → 远程的多级存储保证可靠性。默认5秒将进度持久化到Broker。重启时从Broker加载消费进度继续消费。

2.3、流控管理
为什么要流控? 拉取和处理的异步的,如果处理的速度跟不上拉取速度,会出现内存溢出:缓存无限增长,消费者崩溃:内存耗尽, GC频繁:大量消息在内存中,网络阻塞:大量未处理消息。
RocketMQ消费端用的是队列级流控,据说还有消费者级流控和系统级流控,但是我只找到配置,没找到具体限制的源码。
队列级流控控制的是队列内未处理消息的大小,数量不得超过配置的大小,超过了就一会再过来拉取。
在这里插入图片描述

2.4、保证顺序消费有效性
先看源码,这里说的顺序性是指的顺序消费,位于ConsumeMessageOrderlyService的run,消费的时候会先给对应队列加锁
在这里插入图片描述
再来看看这个锁的获取机制,可以参考这种设计,两级映射结构:MessageQueue → 分片索引 → 锁对象,灵活的锁粒度:支持从整个队列一把锁到多个分片锁,使用 ConcurrentHashMap和putIfAbsent保证线程安全,在保证顺序的同时提高并发度,按需创建锁对象,避免浪费。这边顺序消费传-1保证一个队列只能有一个分片。
在这里插入图片描述

3.消费者负载均衡

在 Apache RocketMQ 领域模型中,同一条消息支持被多个消费者分组订阅,同时,对于每个消费者分组可以初始化多个消费者。以下两种不同的消费效果:
在这里插入图片描述
消费组间广播消费(广播模式) :如上图所示,每个消费者分组只初始化唯一一个消费者,每个消费者可消费到消费者分组内所有的消息,各消费者分组都订阅相同的消息,以此实现单客户端级别的广播一对多推送效果。该方式一般可用于网关推送、配置推送等场景。

消费组内共享消费(集群模式) :如上图所示,每个消费者分组下初始化了多个消费者,这些消费者共同分担消费者分组内的所有消息,实现消费者分组内流量的水平拆分和均衡负载。该方式一般可用于微服务解耦场景。

什么是消费者负载均衡?
如上,广播消费把消息都发给所有消费者了,不涉及负载均衡,这里讨论的是集群共享消费,如果是5.x版本,消费者负载均衡策略分为以下两种模式:

消息粒度负载均衡:PushConsumer和SimpleConsumer默认负载策略
消息粒度负载均衡策略中,同一消费者分组内的多个消费者将按照消息粒度平均分摊主题中的所有消息,即同一个队列中的消息,可被平均分配给多个消费者共同消费。
在这里插入图片描述
如上图所示,消费者分组Group A中有三个消费者A1、A2和A3,这三个消费者将共同消费主题中同一队列Queue1中的多条消息。 注意 消息粒度负载均衡策略保证同一个队列的消息可以被多个消费者共同处理,但是该策略使用的消息分配算法结果是随机的,并不能指定消息被哪一个特定的消费者处理。
消息粒度的负载均衡机制,是基于内部的单条消息确认语义实现的。消费者获取某条消息后,服务端会将该消息加锁,保证这条消息对其他消费者不可见,直到该消息消费成功或消费超时。因此,即使多个消费者同时消费同一队列的消息,服务端也可保证消息不会被多个消费者重复消费。

在顺序消息中,消息的顺序性指的是同一消息组内的多个消息之间的先后顺序。因此,顺序消息场景下,消息粒度负载均衡策略还需要保证同一消息组内的消息,按照服务端存储的先后顺序进行消费。不同消费者处理同一个消息组内的消息时,会严格按照先后顺序锁定消息状态,确保同一消息组的消息串行消费。 顺序消息负载策略
在这里插入图片描述

如上图所述,队列Queue1中有4条顺序消息,这4条消息属于同一消息组G1,存储顺序由M1到M4。在消费过程中,前面的消息M1、M2被消费者Consumer A1处理时,只要消费状态没有提交,消费者A2是无法并行消费后续的M3、M4消息的,必须等前面的消息提交消费状态后才能消费后面的消息。

队列粒度负载均衡:PullConsumer默认负载策略 不做介绍

4.x/3.x版本的负载均衡
队列粒度负载均衡策略中,同一消费者分组内的多个消费者将按照队列粒度消费消息,即每个队列仅被一个消费者消费。这是被大家所熟知的。
在这里插入图片描述
消费端这边提供5种分配策略:
1. AllocateMessageQueueAveragely(默认)平均算法分配,把队列平均分个消费者
2. AllocateMessageQueueAveragelyByCircle 轮询平均分配,就是轮询分配
3. AllocateMessageQueueConsistentHash 一致性哈希分配,使用虚拟节点减少节点变化时的数据迁移,适用于需要最小化重平衡的场景。
4. AllocateMessageQueueByConfig 静态分配,手动指定消费者消费的队列
5. AllocateMessageQueueByMachineRoom 机房就近分配,优先将同机房队列分配给消费者,减少跨机房网络开销。

重平衡机制
RocketMQ 5 的重平衡机制由客户端主动触发,通过while循环,等待20秒去触发重平衡。每个消费者在重平衡时会从 NameServer 获取 Topic 的所有队列和消费者组内所有活跃消费者,然后根据预设的分配策略(如平均分配、一致性哈希等)重新计算自己应该消费的队列,最后通过对比新旧分配结果,释放不再属于自己的队列,并拉取新分配的队列消息,从而实现消费者动态上下线时的负载均衡。整个过程中,消费者之间无需直接通信,通过 NameServer 同步元数据,保证了集群的最终一致性。
在这里插入图片描述

4.消费重试

消费重试指的是,消费者在消费某条消息失败后,Apache RocketMQ 服务端会根据重试策略重新消费该消息,超过一定次数后若还未消费成功,则该消息将不再继续重试,直接被发送到死信队列中。
消息重试的触发条件

消费失败,包括消费者返回消息失败状态标识或抛出非预期异常。
消息处理超时,包括在PushConsumer中排队超时。

消息重试策略主要行为

重试过程状态机:控制消息在重试流程中的状态和变化逻辑。
重试间隔:上一次消费失败或超时后,下次重新尝试消费的间隔时间。
最大重试次数:消息可被重试消费的最大次数。

消息重试策略差异
在这里插入图片描述
这里就只说PushConsumer就好了
在这里插入图片描述
Ready:已就绪状态。消息在Apache RocketMQ服务端已就绪,可以被消费者消费。
Inflight:处理中状态。消息被消费者客户端获取,处于消费中还未返回消费结果的状态。
WaitingRetry:待重试状态,PushConsumer独有的状态。当消费者消息处理失败或消费超时,会触发消费重试逻辑判断。如果当前重试次数未达到最大次数,则该消息变为待重试状态,经过重试间隔后,消息将重新变为已就绪状态可被重新消费。多次重试之间,可通过重试间隔进行延长,防止无效高频的失败。
Commit:提交状态。消费成功的状态,消费者返回成功响应即可结束消息的状态机。
DLQ:死信状态。消费逻辑的最终兜底机制,若消息一直处理失败并不断进行重试,直到超过最大重试次数还未成功,此时消息不会再重试,会被投递至死信队列。您可以通过消费死信队列的消息进行业务恢复。

PushConsumer的最大重试次数由消费者分组创建时的元数据控制,具体参数maxReconsumeTimes。例如,最大重试次数为3次,则该消息最多可被投递4次,1次为原始消息,3次为重试投递次数。超出16次后面都是2小时了。
在这里插入图片描述

5.进度管理

消费进度存储机制
这里分为集群模式和广播模式:
在这里插入图片描述

集群模式采用了"本地缓存+远程持久化"的双层架构。消费者在本地维护进度缓存以提升拉取效率,同时定期将进度同步到Broker端进行集中存储。这种设计的核心价值在于:当消费者重启或发生重平衡时,新分配的消费者能够从Broker获取准确的消费进度,确保消息处理的连续性。由于同组消费者共享同一进度状态,这使得进度集中管理和监控成为可能,为运维提供了统一视图和控制入口。

广播模式则采用了完全分布式的进度管理方案。每个消费者独立消费全量消息,各自维护独立的消费进度,因此无法也没有必要进行集中化管理。若将所有消费者的进度都提交到Broker,会产生大量无关联的分散数据,增加存储和管理的额外开销。基于此,广播模式下每个消费者将进度持久化在本地是最合理的选择,既避免了不必要的资源消耗,也符合其独立的消费语义。

消费点位初始化
消费者可以设置消费点位,具体如下

public enum ConsumeFromWhere {
    CONSUME_FROM_LAST_OFFSET, // 从上次消费位置开始(默认)
    CONSUME_FROM_FIRST_OFFSET, // 从队列头部开始
    CONSUME_FROM_TIMESTAMP; // 从指定时间开始

    private ConsumeFromWhere() {
    }
}

七、RocketMQ-Docker部署(Windows版)

这里以5.3.2为例

1.拉取RocketMQ镜像

docker pull apache/rocketmq:5.3.2

2.创建容器共享网络

RocketMQ 中有多个服务,需要创建多个容器,创建 docker 网络便于容器间相互通信。

docker network create rocketmq

3.启动NameServer

docker run -d --name rmqnamesrv -p 9876:9876 --network rocketmq apache/rocketmq:5.3.2 sh mqnamesrv

验证是否成功,我们可以看到 ‘The Name Server boot success…’, 表示NameServer 已成功启动。

docker logs -f rmqnamesrv

4.启动 Broker+Proxy

# 配置 Broker 的 IP 地址,一定要注意编码字符集问题
echo "brokerIP1=192.168.0.1" > D:\rocketmq\rocketmq-5.3.2\conf\broker.conf

# 启动 Broker 和 Proxy,内存不够的可以自己设置内存
docker run -d 
--name rmqbroker 
--net rocketmq 
-p 10912:10912 -p 10911:10911 -p 10909:10909 
-p 8080:8080 -p 8081:8081 \
-e "NAMESRV_ADDR=rmqnamesrv:9876" 
-e "JAVA_OPT_EXT=-server -Xms1g -Xmx1g -Xmn512m" 
-v D:\rocketmq\rocketmq-5.3.2\conf\broker.conf:/home/rocketmq/rocketmq-5.3.2/conf/broker.conf 
apache/rocketmq:5.3.2 sh mqbroker --enable-proxy \
-c /home/rocketmq/rocketmq-5.3.2/conf/broker.conf

# 验证 Broker 是否启动成功
docker exec -it rmqbroker bash -c "tail -n 10 /home/rocketmq/logs/rocketmqlogs/proxy.log"

我们可以看到 ‘The broker boot success…’, 表示 Broker 已成功启动。

5.部署可视化界面

拉取镜像

docker pull apacherocketmq/rocketmq-dashboard:latest

启动

docker run -d --name rocketmq-dashboard  --net rocketmq  -e "JAVA_OPTS=-Drocketmq.namesrv.addr=rmqnamesrv:9876"  -e "server.port=8082" -p 8082:8082 -t apacherocketmq/rocketmq-dashboard:latest

界面样子
在这里插入图片描述

八、示例代码

我这里用rocketmq-spring-boot-starter简单示例,实际上要自己用原生的rocketmq-client-java,因为spring的跟傻×似的,很多功能不支持(比如批量消费,原生tag设置,配置topic、高级流控、手动ack等等),不过只需要简单的消息发送接收场景可以用,需要复杂灵活场景不行。

1、配置文件

rocketmq:
  name-server: 127.0.0.1:9876
  # 生产者组
  producer:
    group: demo-producer-group
    # 其余走默认配置
    # createTopicKey	发送消息的时候,如果没有找到topic,若想自动创建该topic,需要一个key topic,这个值即是key topic的值	String	TopicValidator.AUTO_CREATE_TOPIC_KEY_TOPIC
    # defaultTopicQueueNums	自动创建topic的话,默认queue数量是多少	int	4
    # sendMsgTimeout	默认的发送超时时间	int	3000,单位毫秒
    # compressMsgBodyOverHowmuc	消息body需要压缩的阈值	int	1024 * 44K
    # retryTimesWhenSendFailed	同步发送失败的话,rocketmq内部重试多少次	int	2
    # retryTimesWhenSendAsyncFailed	异步发送失败的话,rocketmq内部重试多少次	int	2
    # retryAnotherBrokerWhenNotStoreOK	发送的结果如果不是SEND_OK状态,是否当作失败处理而尝试重发	boolean	false
    # maxMessageSize	客户端验证,允许发送的最大消息体大小	int	1024 1024 44M
    # traceDispatcher	异步传输数据接口	TraceDispatcher	null
  # 消费者配置
  consumer:
    # 消费者组
    group: demo-consumer-group
    # 其余走默认配置
    # messageModel	消费模式	MessageModel	MessageModel.CLUSTERINGallocateMessageQueueStrategy	CLUSTERING(集群消費模式) / ROADCASTING (广播消费模式)
    # consumeFromWhere	启动消费点策略	ConsumeFromWhere	ConsumeFromWhere.CONSUME_FROM_LAST_OFFSET
    # consumeTimestamp	CONSUME_FROM_LAST_OFFSET的时候使用,从哪个时间点开始消费	String	半小时前
    # allocateMessageQueueStrategy	负载均衡策略算法	AllocateMessageQueueStrategy	AllocateMessageQueueAveragely(取模平均分配)
    # subscription	订阅关系	Map<String, String>	{}
    # messageListener	消息处理监听器(回调)	MessageListener	null
    # offsetStore	消息消费进度存储器	OffsetStore	null	不建议设置,offsetStore 有两个策略:LocalFileOffsetStoreRemoteBrokerOffsetStore.若沒有显示设置的情況下,广播模式將使用LocalFileOffsetStore,集群模式將使用RemoteBrokerOffsetStore,不建议修改.
    # consumeThreadMin	消费线程池的core size	int	20
    # consumeThreadMax	消费线程池的max size	int	64
    # adjustThreadPoolNumsThreshold	动态扩线程核数的消费堆积阈值	long	100000
    # consumeConcurrentlyMaxSpan	并发消费下,单条consume queue队列允许的最大offset跨度,达到则触发流控	int	2000pullInterval
    # pullThresholdForQueue	consume queue流控的阈值	int	100
    # pullInterval	拉取的间隔	long	0,单位毫秒
    # pullThresholdForTopic	主题级别的流控制阈值	int	-1
    # pullThresholdSizeForTopic	限制主题级别的缓存消息大小	int	-1
    # pullBatchSize	一次最大拉取的批量大小	int	32
    # consumeMessageBatchMaxSize	批量消费的最大消息条数	int	1
    # postSubscriptionWhenPull	每次拉取的时候是否更新订阅关系	boolean	false
    # unitMode	订阅组的单位	boolean	false
    # maxReconsumeTimes	一个消息如果消费失败的话,最多重新消费多少次才投递到死信队列	int	-1	由于PullConsumer没有管理消费的线程池和管理器,需要用户自己处理各种消费结果和拉取结果,故需要投递到重试队列或死信队列的时候需要显示调用sendMessageBack.回传消息的时候会带上maxReconsumeTimes的值,broker发现此消息已经消费超过此值,则投递到死信队列,否则投递到重试队列。此逻辑和DefaultPushConsumer是一致的,只是PushConsumer无需用户显示调用.
    # suspendCurrentQueueTimeMillis	串行消费使用,如果返回ROLLBACK或者SUSPEND_CURRENT_QUEUE_A_MOMENT,再次消费的时间间隔	long	1000
    # consumeTimeout	消费的最长超时时间	long	15,单位分钟
    # awaitTerminationMillisWhenShutdown	关闭使用者时等待消息的最长时间,0表示无等待。	long	0
    # traceDispatcher	异步传输数据接口	TraceDispatcher	null
    # registerTopics	消費者需要監聽的topic	Collection	默認值:空集合

2、生产者示例

单向发送,没有返回结果。

  /**
     * 发送单向消息
     *
     * @param msg 消息
     * @return 响应
     */
    @GetMapping("/sendOneWay")
    public String sendOneWay(String msg) {
        Message<String> message = MessageBuilder.withPayload(msg).build();
        rocketMQTemplate.sendOneWay("demo-topic", message);
        return "success";
    }

同步发送,阻塞等待,有返回结果。

/**
     * 发送同步消息
     *
     * @param msg 消息
     * @return 响应
     */
    @GetMapping("/sendSync")
    public SendResult sendSync(String msg) {
        Message<String> message = MessageBuilder.withPayload(msg).build();
        return rocketMQTemplate.syncSend("demo-topic", message);
    }

异步发送,不需要阻塞,回调通知。

/**
     * 发送异步消息
     *
     * @param msg 消息
     * @return 响应
     */
    @GetMapping("/sendAsync")
    public void sendAsync(String msg) {
        Message<String> message = MessageBuilder.withPayload(msg).build();
        rocketMQTemplate.asyncSend("demo-topic", message, new SendCallback() {
            @Override
            public void onSuccess(SendResult sendResult) {
                System.out.println("发送成功:" + sendResult);
            }

            @Override
            public void onException(Throwable throwable) {
                System.out.println(throwable.getMessage());
            }
        });
    }

顺序消息,根据队列策略计算选择的队列,只有都往一个队列发送才能保证顺序性。

 /**
     * 顺序发送消息
     *
     * @param msg 顺序消息
     * @return 响应
     */
    @GetMapping("/sendOrderly")
    public SendResult sendOrderly(String msg) {
        Message<String> message = MessageBuilder.withPayload(msg).build();
        return rocketMQTemplate.syncSendOrderly("demo-topic-orderly", message, "orderly");
    }

延迟消息,按照延迟级别,到时间投递。

  /**
     * 延时发送消息
     *
     * @param msg 延时消息
     * @return 响应
     */
    @GetMapping("/sendDelay")
    public SendResult sendDelay(String msg, int delayLevel) {
        Message<String> message = MessageBuilder.withPayload(msg).build();
        // delayLevel: 1s 5s 10s 30s 1m 2m 3m 4m 5m 6m 7m 8m 9m 10m 20m 30m 1h 2h
        return rocketMQTemplate.syncSend("demo-topic", message, 3000, delayLevel);
    }

带tag和key的消息发送,tag一般用来消费者过滤想要的消息,key一般用来标识消息唯一性,比如主键id,或者用于消息跟踪,tag要拼接到topic后面。

/**
     * 发送同步消息带tag和key
     *
     * @param msg 消息
     * @return 响应
     */
    @GetMapping("/sendSyncWithTag")
    public SendResult sendSync(String msg, String tag, String key) {
        Message<String> message = MessageBuilder
                .withPayload(msg)
                .setHeader(RocketMQHeaders.KEYS, key)
                .build();
        return rocketMQTemplate.syncSend("demo-topic:" + tag, message);
    }

批量发送。

  /**
     * 批量发送消息
     *
     * @param msgs 批量消息
     * @return 响应
     */
    @PostMapping("/sendBatch")
    public SendResult sendBatch(@RequestBody List<String> msgs) {
        Collection<Message<?>> messageCollection = new ArrayList<>();
        for (String msg : msgs) {
            Message<String> message = MessageBuilder.withPayload(msg).build();
            messageCollection.add(message);
        }
        return rocketMQTemplate.syncSend("demo-topic-batch", messageCollection);
    }

发送事务消息,可以在head和方法里传参,建议在head传,因为事务回查的时候没有object参数,看下面。

 /**
     * 事务消息发送
     *
     * @param msg 消息
     * @return 响应
     */
    @GetMapping("/sendTransaction")
    public TransactionSendResult sendTransaction(String id,String msg) {
        Message<String> message = MessageBuilder
                .withPayload(msg)
                .setHeader("test_id", id)
                .build();
        return rocketMQTemplate.sendMessageInTransaction("demo-topic", message, "参数");
    }

配置事务事件监听器,用于事务回查和本地事务提交,我这里模拟id1事务提交成功,消息会真正发出去,id2提交失败,消息不会发出去,消息3等待3次才真正提交成功,才会发出去。
RocketMQLocalTransactionState.COMMIT 事务提交
RocketMQLocalTransactionState.ROLLBACK 事务回滚
RocketMQLocalTransactionState.UNKNOWN 暂不提交,会执行checkLocalTransaction,不过有个退避策略,最大回查15次就当做失败了,每次回查时间会递增。

/**
 * @author chendx
 * @date 2026/1/16 17:14
 * @description: 事务监听器
 */
@Component
@RocketMQTransactionListener
public class RocketMQTranscationListener implements RocketMQLocalTransactionListener {
    private AtomicInteger count = new AtomicInteger(0);
    @Override
    public RocketMQLocalTransactionState executeLocalTransaction(Message message, Object o) {
        System.out.println("执行本地事务");
        System.out.println(message);
        System.out.println(o);
        String id = message.getHeaders().get("test_id").toString();
        if ("1".equals(id)) {
            // 提交事务
            return RocketMQLocalTransactionState.COMMIT;
        } else if ("2".equals(id)) {
            // 回滚事务
            return RocketMQLocalTransactionState.ROLLBACK;
        } else if ("3".equals(id)) {
            // 未知状态,继续查询
            return RocketMQLocalTransactionState.UNKNOWN;
        }
        return null;
    }

    @Override
    public RocketMQLocalTransactionState checkLocalTransaction(Message message) {
        System.out.println("查询本地事务");
        String id = message.getHeaders().get("test_id").toString();
        if ("1".equals(id)) {
            return RocketMQLocalTransactionState.COMMIT;
        } else if ("2".equals(id)) {
            return RocketMQLocalTransactionState.ROLLBACK;
        } else if ("3".equals(id)) {
            int i = count.addAndGet(1);
            if (i < 3) {
                return RocketMQLocalTransactionState.UNKNOWN;
            }
            return RocketMQLocalTransactionState.COMMIT;
        }
        return null;
    }
}

3、消费者示例


最简单的一个消费者。

/**
 * @author chendx
 * @date 2026/1/16 9:22
 * @description: 默认就是并发集群消费,20个线程,最大64个
 */
@Component
@RocketMQMessageListener(
        consumerGroup = "demo-group",
        topic = "demo-topic")
public class RocketMQConsumerDemo implements RocketMQListener<String> {

    @Override
    public void onMessage(String msg) {
        System.out.println("线程" + Thread.currentThread().getName() + ",默认监听器收到消息:" + msg);
    }
}

带tag过滤的消费者。

/**
 * @author chendx
 * @date 2026/1/16 9:22
 * @description: 过滤tag1和tag2的
 */
@Component
@RocketMQMessageListener(
        consumerGroup = "demo-tag-group",
        topic = "demo-topic",
        selectorExpression = "tag1 || tag2")
public class RocketMQConsumerTagDemo implements RocketMQListener<String> {

    @Override
    public void onMessage(String msg) {
        System.out.println("线程" + Thread.currentThread().getName() + ",tag监听器收到消息:" + msg);
    }
}

顺序消息消费者,关键是consumeMode = ConsumeMode.ORDERLY

/**
 * @author chendx
 * @date 2026/1/16 9:22
 * @description: 顺序消费
 */
@Component
@RocketMQMessageListener(
        consumerGroup = "demo-orderly-group",
        topic = "demo-topic-orderly",
        consumeMode = ConsumeMode.ORDERLY)
public class RocketMQConsumerOrderlyDemo implements RocketMQListener<String> {

    @Override
    public void onMessage(String msg) {
        System.out.println("线程" + Thread.currentThread().getName() + ",顺序监听器收到消息:" + msg);
    }
}

总结

以上就是今天要讲的内容,本文对RocketMQ的核心概念、架构部署、实践应用和高级特性进行了讲解。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐