RocketMQ之初识MQ
文章目录
前言
随着微服务架构的普及以及大数据时代的到来,传统的单体服务已经无法满足“高并发”、“高可用”的需求,微服务架构采用多实例部署分均服务压力,也提升了服务的可用性,但微服务随之而来的问题则是如何保证数据的一致性;服务之间通信如何保证可靠;突发流量时,如何做到流量控制,而MQ的出现正好可以解决以上的问题。
一、MQ是什么?
MQ是Message Queue的缩写,我们一般称为消息队列,基于数据结构Queue(队列)实现,队列遵循先进先出(FIFO)的原则,而MQ基本也是遵循这个原则,只是在某些延迟队列时,是按照设置的延迟时间排序进行消费。MQ基于观察者模式实现发布-订阅的机制,观察者模式定义对象间的“一对多”依赖关系,当被观察对象状态变化时,所有观察者自动收到通知并更新。MQ通过生产者发布消息、消费者订阅消息的方式实现异步通信,本质上也是一对多的消息分发机制。
总得来说消息队列(Message Queue,简称MQ)是一种在应用程序之间传递消息的通信方法。它允许应用程序通过发送和接收消息来进行异步通信,从而解耦系统组件,提高系统的可扩展性和可靠性。
MQ一般由三大核心概念组成:
- 生产者(Producer):负责产生消息并发送到消息队列。
- 消费者(Consumer):从消息队列中获取消息并进行处理。
- 消息队列(Message Queue)用于接收生产者发送的消息;存储消息避免消息丢失;转发消息给对应的消费者。
MQ的消费模式一般有两种组成:
- 点对点模式(Point-to-Point):生产者将消息发送到队列,消费者从队列中获取消息。一条消息只能被一个消费者消费。
- 发布/订阅模式(Publish/Subscribe):生产者将消息发送到主题(Topic),多个消费者可以订阅这个主题,每个消费者都会收到消息的副本。
MQ的主要特点有以下几点:
- 异步通信:生产者发送消息后无需等待消费者立即处理,可以继续执行其他任务。
- 解耦:生产者和消费者之间不直接相互调用,降低了系统间的依赖。
- 削峰填谷:在流量激增时,消息队列可以缓存消息,避免系统被冲垮,然后消费者按照自己的处理能力消费消息。
- 可靠性:许多消息队列提供持久化机制,确保消息不会丢失。
- 扩展性:由于解耦,可以独立地扩展生产者和消费者。
常见消息队列产品有以下四个:
| 名称 | 功能支持 | 时效性 | 单机吞吐量 | Client SDK |
|---|---|---|---|---|
| Kafka | 分布式、分区、高吞吐量、流处理、持久化 | 毫秒级 | 高(百万级消息每秒) | 多语言(JAVA\Scala\Python) |
| ActiveMQ | JMS 支持、持久化、事务、消息过滤 | 毫秒级 | 中等(十万级消息每秒) | 多语言(JAVA\C\C++) |
| RabbitMQ | 消息路由、插件机制、持久化、多种协议支持 | 微秒级到毫秒级 | 中等(十万级消息每秒) | 多语言(JAVA\Erlang.NET) |
| RocketMQ | 分布式、高吞吐量、事务、定时和延时消息、持久化 | 毫秒级 | 高(百万级消息每秒) | 多语言(JAVA\C++\GO) |
二、RocketMQ 是什么?
RocketMQ 是阿里开源的消息中间件,具有高性能、高可靠、高实时、分布式 的特点,底层是用 Java 语言开发的分布式组件。Apache RocketMQ 自诞生以来,因其架构简单、业务功能丰富、具备极强可扩展性等特点被众多企业开发者以及云厂商广泛采用。历经十余年的大规模场景打磨,RocketMQ 已经成为业内共识的金融级可靠业务消息首选方案,被广泛应用于互联网、大数据、移动互联网、物联网等领域的业务场景。官方文档:RocketMQ官网,源码地址:rocketmq。以下展示的是RocketMQ的基本架构图:

可以看到RocketMQ主要分为四个部分,分别是Producer、Consumer、NameServer、Broker。
(2.1) NameServer—名字服务器
NameServer是RocketMQ中的核心组件,是RocketMQ中的注册中心,支持Broker的动态注册和发现。主要包括两个功能:
- Broker管理
NameServer管理Broker的注册信息,包括Broker的IP地址,Topic信息。将这些信息做为路由信息的基本数据,提供心跳检测机制,检查Broker是否存活。
- 路由信息管理
每个NameServer将保存关于 Broker 集群的整个路由信息和用于客户端查询的队列信息。Producer和Consumer通过NameServer可以知道整个Broker集群的路由信息,从而进行消息的投递和消费。
NamerServer默认是支持集群多实例的,各个实例之间是不进行信息通讯的,因此NameServer之间几乎是无状态的。当启动Broker时,需要填写多个NameServer的地址(如果NameServer有多台),以至于Broker向每台NameServer发送注册自己的路由信息,每个NameServer实例上都会保存一份完整的路由信息,当某个NameServer因某种原因下线了,客户端仍然可以向其它NameServer获取路由信息。

(2.2) Broker—代理服务器
Broker 主要负责消息的存储、投递和查询以及服务高可用保证,处理生产者和消费者的消息读写请求。Broker在启动时,会向所有的NameServer注册自己的信息(包括IP地址、Topic信息等),并定期发送心跳以保持连接,Broker主要提供以下功能:
- 消息存储:Broker将消息持久化到磁盘(通常采用顺序写文件的方式,以保证高性能)。
- 消息投递:Broker根据消费者的请求(拉取或推送)将消息发送给消费者。
- 高可用:Broker可以配置为主从结构,主Broker负责处理读写请求,从Broker负责从主Broker同步数据,以实现高可用。
在 Master-Slave 架构中,Broker 分为 Master 与 Slave。一个Master可以对应多个Slave,但是一个Slave只能对应一个Master。Master 与 Slave 的对 应关系通过指定相同的BrokerName,不同的BrokerId 来定义,BrokerId为0表示Master,非0表示Slave。Master也可以部署多个。
(2.3) Producer—生产者
生产者在RocketMQ架构中,作为用来构建并传输消息到服务器(Broker)的运行实体。通过不断的与NameServer集群中某一个NameServer进行通信,拉取最新的路由信息,在发送消息时,根据拉取的路由信息以及对应的策略,将消息发送到某个Broker上进行消息存储、投递。
在消息生产者中,可以定义如下传输行为:
- 发送方式:生产者可通过API接口设置消息发送的方式。Apache RocketMQ 支持同步传输和异步传输。
- 批量发送:生产者可通过API接口设置消息批量传输的方式。例如,批量发送的消息条数或消息大小。

(2.4) Consumer—消费者
消费者是RocketMQ 中用来接收并处理消息的运行实体,消息者启动时,需要配置其通信的NameServer集群的IP,以便于与NameServer集群中的某个NameServer间建立通信,获取Broker集群的路由信息。消费者消费消息通常有两种模式:推(Push)和拉(Pull)。
| 特性 | push模式 | pull模式 |
|---|---|---|
| 消息获取 | Broker主动发起 | Consumer主动拉取 |
| 实时性 | 高(消息到达立即推送) | 依赖拉取间隔 |
| 流量控制 | 消费者通过暂停请求控制 | 消费者直接控制拉取频率 |
| 资源消耗 | Broker需要维护推送连接 | 消费者承担轮询开销 |
通过以上表格我们知道Pull模式由消费者自主发起,因此一般不存在自身消费压力过大的情况。而Push模式是由Broker发起,但消费者这边有控制权,不会出现无脑推送消息导致是消费者消费压力过大而挂掉的情况。

(2.5) 小结总结
- 每一个Broker都会与NameServer集群中的所有节点建立长连接,定时向每个节点注册自身的信息。
- 一个Producer只会与NameServer集群中的一个节点建立长连接,定期从NameServer获取Topic路由信息,并向提供 Topic 服务的 Master (Broker)建立长连接,且定时向 Master (Broker)发送心跳。
- 一个Consumer只会与 NameServer 集群中的其中一个节点建立长连接,定期从 NameServer 获取 Topic 路由信息,并向提供 Topic 服务的 Master、Slave 建立长连接,且定时向Master、Slave发送心跳。Consumer 既可以从 Master 订阅消息,也可以从Slave订阅消息。
生产者只与Broker集群中的Master建立连接通信,不会对Slave进行连接,Salve只针对Master做数据同步。在读写分离的情况下,Master处理写请求(生产消息),Slave处理读请求(消费消息),两者通过主从复制保持数据一致性。
消费者可以与Broker集群中的Master和Slave建立连接,当Master节点因故障无法处理消息时,Consumer会自动切换到Slave节点获取消息,这个过程无需人工干预。
三、RocketMQ 主题Topic
在介绍了RocketMQ的基础概念(名字服务器、代理服务器、生产者、消费者)后,接下来介绍消费主体-Topic
在官方文档中,将Topic定义为消息传输和存储的顶级容器,用于标识同一类业务逻辑的消息。并且建议将不同业务类型的消息拆分到不同的主题Topic中进行管理,通过主题能够实现存储的隔离性和订阅的隔离性。每条消息必须指定属于某个主题,要注意主题是一个逻辑概念,并不是实际存储消息的容器。一个主题是由一个或者多个队列组成,队列才是存储消息的实际容器,消息的存储和水平扩展能力最终是由队列实现的;并且针对主题的所有约束和属性设置,最终也是通过主题内部的队列来实现。

Topic 主要有 2 大核心作用:
- 定义数据的分类隔离:将不同业务类型的数据拆分到不同的主题中管理,通过主题实现存储的隔离性和订阅隔离性
- 定义数据的身份和权限:RocketMQ 的消息本身是匿名无身份的,同一分类的消息使用相同的主题来做身份识别和权限管理。
(3.1) Topic 相关概念与配置
(3.1.1) 主题命名规则
主题的名称是在创建主题的时候需要我们自定义,名称最好与业务类型相关,名称需要满足在整个集群内唯一即可。Topic命名应该尽量使用简短、常用的字符,避免使用特殊字符。特殊字符会导致系统解析出现异常,字符过长可能会导致消息收发被拒绝。

(3.1.2) 队列(MessageQueue)
队列是RocketMQ中消息存储和传输的实际容器,也是RocketMQ中消息的最小存储单元,所有主题都是由至少一个队列组成,所有发送成功的消息都被持久化存储到队列中,配合生产者和消费者客户端的调用可实现至少投递一次的可靠性语义,通过修改队列数量,以此实现横向的水平扩容和缩容。

队列的主要作用有如下两点:
存储顺序性
队列天然具备顺序性,即消息按照进入队列的顺序写入存储,同一队列间的消息天然存在顺序关系,队列头部为最早写入的消息,队列尾部为最新写入的消息。消息在队列中的位置和消息之间的顺序通过位点(Offset)进行标记管理。
流式操作语义
Apache RocketMQ 基于队列的存储模型可确保消息从任意位点读取任意数量的消息,以此实现类似聚合读取、回溯读取等特性,这些特性是RabbitMQ、ActiveMQ等非队列存储模型不具备的。
队列数可在创建主题或变更主题时设置修改,队列数量的设置应遵循少用够用原则,避免随意增加队列数量。
消费者从订阅的主题,对应主题下的队列获取消息进行消费,当发现队列消息消息堆积时,可以适当提升消费者数量增加消费速度。不过应该保持消费者数量<=队列数量。在集群模式下,一个队列下同一个消费者组内,只能被一个消费者消费。一个队列只能被一个消费者绑定,如果消费者数量大于队列数量,那么多余的消费者将无法消费消息,所以,在设计时,通常要让消费者组内的消费者数量小于等于队列数量,以避免资源浪费。
(3.1.3) 消息类型
消息类型是RocketMQ主题所支持的消息类型,在创建主题时需要选择消息类型,主要包括Normal(普通消息)、FIFO(顺序消息)、Delay(定时/延时消息)、Transaction(事务消息)。
- Normal:
是普通消息,是 RocketMQ默认的消息类型,消息本身无特殊语义,消息之间也没有任何关联。普通消息没有复杂的处理逻辑,适用于大多数常见的消息传递场景
- FIFO:
顺序消息,消息按照发送的顺序被严格地消费,适用于对消息顺序有严格要求的场景,比如数据实时增量同步。RocketMQ
支持两种类型的顺序消息:全局顺序消息和分区顺序消息。 全局顺序消息保证同一个 Topic下的所有消息按照严格的顺序被消费。这意味着所有消息按照它们被发送的顺序依次被消费。这种模式比较少用,因为它限制了消费并发性。
分区顺序消息保证某一个分区(即某一个队列)内的消息按顺序消费,但不同分区之间的消息消费顺序不做保证。分区顺序消息通过将消息分配到不同的队列,实现了更高的并发性能。
- Delay:
定时/延时消息,通过指定延时时间控制消息生产后不要立即投递,而是在延时间隔后才对消费者可见。这种消息类型很有意思,适用于延迟任务、任务超时处理、分布式定时调度等场景。这种场景我们后面会详细介绍。
- Transaction:
这种消息类型为事务消息,也是基于消息队列来实现分布式事务,持应用数据库更新和消息调用的事务一致性保障。
Rocket5.0开始,支持强制校验信息类型,即每个主题只允许发送一种消息类型的消息,不允许将多种类型的消息发送到同一个主题中。
(3.1.4) Topic 主题总结
RocketMQ默认提供自动创建主题的功能,在生产者发送消息时,如果发现主题不存在,则Broker会检查autoCreateTopicEnable配置,这个值默认时True,表示当主题不存在时会自动创建主题。由于主题属于顶层资源和容器,应该拥有独立的权限管理、可观测性指标采集和监控等能力,创建和管理主题会占用一定的系统资源。因此,生产环境需要严格管理主题资源,建议将这个值设置为False。
Topic是RocketMQ中概念模型的核心部分,用于标识同一类业务逻辑的消息,主题拆分既不能太细又不能太粗,一般可以依据以下角度考虑拆分粒度:
- 消息类型是否一致
不同类型的消息,如顺序消息和普通消息需要使用不同的主题
-消息业务是否关联
如果业务没有直接关联,微信小程序订阅消息和微信群消息接收两个业务是不相关的,所以主题也肯定应该是分开的。
- 消息量级别是否一致
数量级别或者时效性不同的业务消息应该使用不同的主题,比如消息量小但是时效性要求高的消息与消息量大的消息应该划分到不同的主题中去。
四、消息队列基础模型
在消息队列中,一般秉承着两种基础模型,一般称为队列模型(点对点)和主题模型(发布/订阅)。
- 队列模型:
其原理基于数据结构中的队列,核心思想为竞争消费,一条消息只会被一个消费者处理,点到点的消费模式。
如果多个消费者同时监听一个队列,消费队列系统会确保一个条信息只会被其中一个消费者获取和消费,这自然地在多个消费者之间形成了负载均衡,提高了系统的处理能力。一旦消费成功,就从队列中移除。比如银行柜台叫号系统等。
- 主题模型
主题模型的核心思想是广播消息,一条信息可被多个消费者处理。即生产者将消息发送到主题后,只有订阅了该主题的消费者才有可能消费到该条消息。其特性为一对多,即发送到主题的一条消息,会被所有订阅了该主题的消费者各自收到一份副本。常用于消息订阅,比如微信公众号订阅,更新会发送消息给你。当一个事件发生时(如“用户注册成功”),需要通知多个独立的下游系统。用户服务发布一条“用户注册成功”的消息到主题,积分服务、邮件服务、推荐服务等各自订阅该主题并执行相应的逻辑。
现代主流的消息队列(比如RocketMQ)不再是单纯的主题模型或者是队列模型,是通过"主题+消费组"的机制将两者完美融合。
- 主题:仍然是逻辑上的消息分类,用于实现发布/订阅的广播特性。
- 消费者组:由多个消费者实例组成,这些实例功能完全相同,目的是为了负载均衡和容错。
在这种模型下,产生了两种主要的消费模式:
- 集群模式 - 实现了队列模型的效果
规则:同一个消费者组内的所有消费者,以竞争的方式消费主题下的消息(
一条消息只会被组内的一个消费者消息)
效果:在消费组内部,实现了点对点的负载均衡
目的:用于横向扩展消费者的处理能力,避免重复消费
- 广播模式 - 实现了纯粹主题模型的效果
规则:同一个消费组的所有消费者,每个人都能收到主题的全部消息。
效果:消费组内部实现发布/订阅的广播
目的:用于确保每个消费者实例都拥有全量数据
广播模式是订阅了该主题的所有消费者都能够消费同一个条消息,而集群模式下,是按照消费组来进行消息,一个主题下的消息,可以被多个消费组进行消费,但是消费组存在多个消费者时,对这条消息是存在竞争关系,即一个消费组下,只有一个消费者能够消费该条消息。
(4.1) 消息点位(Message Offset)和消费点位(Consumption Offset)
(4.1.1)消息点位(Message Offset)
主题Topic只是一个逻辑概念,其由多个消息队列MessageQueue,消息实际存储在文件中,每个消息都有一个唯一的坐标,有点类似于数组下标用于定位消息在队列中的位置.。消息位点是消息在消息队列存储文件中的物理偏移地址,这个坐标称为消息点位**(Message Offset)**,也称为消息偏移量。消息位点用于唯一标识消息,并记录消息在队列中的位置。消费者通过位点来记录自己消费到了哪里,以便后续从该位置继续消费。

在消息队列 (MessageQueue) 中最早的一条消息叫 最小消费位点(MinOffset),最新的一条消息就叫 最大消费位点。理论上最小和最大消费位点之间可以被无限拉大,但是受物理服务器内存存储限制,一般消息文件会被保留72小时,可以通过 Broker 配置文件中的 fileReservedTime 参数来控制消息文件的存储时间,单位为小时。当某个消息文件被删除后,最小消费位点就会更新为当前有效消息文件中的队列的最早消息的存储位置。消息位点有以下特性:
- Broker分配
消息位点是消息在消息队列存储文件中的物理偏移地址,由消息队列系统(Broker)在消息持久化时自动分配,生产者或消费者无法修改。
- 唯一性
在同一个队列(Queue)或分区(Partition)内,每个消息位点都是唯一的
- 顺序性
位点的大小决定了消息被存储和投递的先后顺序。位点小的消息一定先于位点大的消息被存储
- 持久化
消息和它的位点一起被持久化到磁盘,不会丢失
(4.1.1)消费点位(Consumption Offset)
消费位点又称为消费进度或者消费偏移量,消费位点记录了某个 消费者组(Consumer Group) 对某个 队列(Queue) 的消费进度。它表示"这个消费者组已经消费到了这个位置,之后的消息还没有消费。其存在以下特性:
- 消费者组绑定
消费位点是和消费者组关联的,而不是和单个消费者客户端关联。同一个主题的一条消息,不同的消费者组会有各自独立的消费位点。
- 可前进
随着消费者不断消费和提交位点,消费位点会单调向前推进
- 需要持久化
在特定情况下(如需要重新消费历史消息),消费位点可以被重置到一个更早的位置
- 需要持久化
消费位点必须被可靠地存储起来,以便消费者重启后能从上次的位置继续消费,避免消息重复或丢失
消费位点提交策略直接决定消息的消费语义,通常由"至少消费一次"、“至多一次”、“恰好一次”。
- 至少一次
处理流程:费者在处理完消息之后再提交消费位点
风险:如果处理成功后,在提交位点之前消费者崩溃,新的消费者会从旧位点开始重新消费,导致消息重复
保存:消息绝不会丢失
- 至多一次
处理流程:消费者在处理消息之前就先提交消费位点
风险:如果提交位点后,处理消息之前消费者崩溃,消息将不会被处理,导致消息丢失
保存:消息绝不会重复
- 恰好一次
处理流程:这是最理想但最难实现的。通常需要外部事务(如数据库事务)的配合,将消息处理和位点提交放在同一个事务中,保证其原子性。或者在流处理引擎中通过幂等和事务来实现
作用:它是消费进度的"指针"或"书签",确保了消息"至少被处理一次"或"恰好被处理一次"的语义,类似于多个读者(消费组)共同阅读一本书,每个消费组都可以单独设置书签(消费位点)
(4.2) 消息堆积处理
(4.2.1)消息堆积处理
消息堆积顾名思义就是消息大量堆积在消息队列中,无法快速被消费掉,造成消息堆积主要的原因有如下几点:
- 消费者处理能力不足
消费者处理能力不足一般可能体现在消费者数量过少、单个消费者消费处理速度过慢或者消费者出现异常宕机情况等。
- 生产者突发流量
生产者突然流量一般指的是短时间内大量消息涌入,导致生产者产生消息速率过快。
- 网络延迟或者不稳定
网络延迟或者不稳定是比较难排查的问题,其影响最大的是可能阻塞消费者与broker之间的通信,导致在消息消费时间大部分花费在通信上,降低了消费者消费消息的速度。
- 硬件资源不足
当消费者所在的服务器CPU、内存等资源不足时,通常会造成磁盘I/O卡顿,从而降低消费者消费速度。
处理消息堆积的方法有很多,通常就是优化消费者消费逻辑,提高单个消费者处理速度。或者就是调整消费者配置,增加消费者数量,增加主题消息队列、提高硬件配置等待。最常见解决消费堆积的方式就是增加消费者和增加队列,不过并不是队列越多越好,也不是消费者数量越多越好。
(4.2.1.1)消费并行度概念
在RocketMQ 的消费模型中,一个队列在同一个时刻只能被一个消费者线程消费。也就是说一个队列同一个时刻只能绑定一个消费者。因此,整个消费集群的最大并行度受限于两者中的较小值。
- 有 8 个队列,10个消费者,由于一个队列只会被一个消费同时绑定,因此最多同时有8个消费者能够进行消费,也就是最大并行度为8,而剩余的2个消费者将会空闲。
- 有 10 个队列,5个消费者,那么每个消费者可以分配到2个队列进行消费,同时最多5个消费者可以同时进行消费,因此最大并行度为5。
得出结论:在不超过系统瓶颈的前提下,为了让并行度最大化,我们通常会让消费者数量 ≈ 队列数量,因此消费者数量并不是越多越好
(4.2.1.2)队列数量过多的影响
队列是RocketMQ进行消息负载和存储的基本单位,盲目的增加队列数量会带来以下影响:
- Broker端性能与资源消耗
(1)内存压力:每个队列在 Broker 上都会维护一个独立的消费位点(consumer offset)。队列数量越多,需要管理的元数据就越多,消耗的内存也越大
(2)文件句柄与IO压力:每个队列对应存储层面的多个文件(CommitLog 和 ConsumeQueue)。队列数激增会导致文件数量暴涨,消耗大量文件句柄,在消息刷盘和同步时可能引发 IO 竞争,反而降低整体吞吐量。
(3)HA 与同步性能:在主从同步模式下,每个队列都需要进行主从状态同步。队列数量过多会加重主从同步的负担,增加网络开销,在故障切换时也可能导致更长的恢复时间。
- NameServer 与运维复杂度
元数据膨胀:topic 的路由信息(包含所有队列信息)需要存储在 NameServer 并推送给所有客户端。队列数量过多会导致路由数据变得臃肿,增加网络传输和内存开销。
运维监控困难:在管理控制台(RocketMQ Console)上监控成千上万个队列的状态是非常困难的,告警和问题排查的粒度也会变得很粗。
- 客户端性能受限
客户端负载:Producer 和 Consumer 在启动时需要从 NameServer 拉取完整的路由信息。队列数量过多,这个路由表就会很大。
Consumer 端:Rebalance(负载重均衡)过程会变得更复杂、更耗时。当消费者组内的消费者数量发生变化时,需要重新分配大量的队列,在此期间消费会暂停。
(4.3) 消息的过滤方式
在讨论RocketMQ消息过滤方式之前,我们应该先讨论为何要有消息过滤这个概念。根据上面小节我们知道一个主题Topic对应了一个或者多个队列,而订阅者(消费者)通过订阅主题来进行消息消费。但是主题是一个较大的逻辑概念,比如我们可以把短信发送信息定义为一个主题。但是可以试想一下,在发布订阅模型下,RocketMQ 会将所有订阅了主题的消息都投递给消费者,但有时候消费者只关⼼消息⾥的某⼀类型的内容⽽不是全量消息。比如短信信息中又可以分为登录验证码信息、注册验证码信息、用户订阅某某某产品信息、密码过期提示短信信息等。而按照单一原则,应该由不同的消费者来处理这个不同的信息。用户订阅产品信息的消费者不应该关心登录验证码信息、注册验证码信息等。因此在主题Topic中再次进行数据过滤是十分有必要的!

(4.3.1)RocketMQ如何进行消息过滤
在RocketMQ中有一个强大的功能,就是可以对一条信息进行打上标签,或者对消息进行自定义属性,从而达到消息识别的效果。⽣产者在发送消息时,可以对即将发送的消息进行打上标签标识,而在消费端消费消息时,可以指定从订阅的主题中只拿取某一种标签标识的消息进行消费,即:
- 生产者发送消息时,给消息定义属性或者标签
- 消费者在调用主题订阅注册接口时,告知RocketMQ自身只想消费带有哪些标签或者属性的消息
- RocketMQ 服务端:消费者获取消息时,会触发服务端的过滤算法,按需给消费者对应消息消息
在RocketMQ中消息过滤分为两大类型,一种是Tag标签过滤,一种是SQL属性过滤。Tag消息过滤是最常见的消息过滤类型,在生产者产生消息标签时,对消息设置对应标签即可,实现代码如下:
public class RocketMQMessage {
public static void main(String[] args) {
Message message = new Message();
message.setTopic("订阅的主题");
message.setTags("TagA");//设置消息Tag
message.setBody("消息体".getBytes());
}
}
在消费端设置消费者指定消费某个标签的消息,代码如下:
DefaultMQPushConsumer defaultMQPushConsumer = new DefaultMQPushConsumer(singleConsumer.getGroup());
defaultMQPushConsumer.setNamesrvAddr("127.0.0.1");
defaultMQPushConsumer.setVipChannelEnabled(false);
defaultMQPushConsumer.setInstanceName("Instance_Name");
defaultMQPushConsumer.setConsumeFromWhere(ConsumeFromWhere.CONSUME_FROM_FIRST_OFFSET);
defaultMQPushConsumer.setConsumeMessageBatchMaxSize(1);
defaultMQPushConsumer.setMessageModel(MessageModel.CLUSTERING);
defaultMQPushConsumer.setAllocateMessageQueueStrategy(new AllocateMessageQueueAveragely());
defaultMQPushConsumer.setMaxReconsumeTimes(16);
defaultMQPushConsumer.setConsumeTimeout(10);
defaultMQPushConsumer.setMaxReconsumeTimes(16);
defaultMQPushConsumer.setConsumeThreadMin(16);
defaultMQPushConsumer.setConsumeThreadMax(8);
defaultMQPushConsumer.subscribe('订阅的主题', 'TagA');//设置消费者订阅的主题和标签
五、RocketMQ的消息重试机制
消息重试的⽬的是为了保证消息的完整性,防⽌消息丢失,是业务兜底策略。消息重试一般发生在以下一些场景:
- 内部数据自己处理异常,也就是自身系统有编码功能缺陷而导致发送或者接收消息后处理异常。
- 网络故障:通常是网络不可达或者网络丢包发生,需要尝试重新发送信息
- 消费者处理消息时突然宕机,可能是服务升级或者服务内存溢出等问题导致
消息重试机制能解决以下问题:
- 临时性故障处理: 当消费者因为⽹络波动、服务器暂时不可⽤等临时性问题⽆法正常消费消息时,重试机制可以在稍后再次尝试投递,提⾼消息处理的成功率。
- 业务处理异常: 如果消费者在处理消息时遇到业务逻辑异常,重试机制可以让消费者有机会在稍后再次处理该消息。
- 提⾼系统可靠性: 通过多次重试,可以降低因单次失败导致的消息丢失⻛险,提⾼整个消息系统的可靠性和稳定性。
- 应对突发流量: 在消费者遇到突发⾼流量⽆法及时处理所有消息时,重试机制可以帮助削峰填⾕,让消息在稍后的低峰期被重新处理。
- 分布式事务处理: 在实现最终⼀致性的分布式事务中,消息重试可以作为⼀种补偿机制,确保事务的最终完成。

重试⼜分⽣产端的消息重试和消费端的消息重试
(5.1) ⽣产端的消息重试
⽣产者向 RocketMQ 的 broker 发送消息时,因⾃身原因如⽹络波动等导致消息没有投递成功,此时就需要进行消息重发。在生产端,不同的发送方式有着不同的重试策略,RocketMQ发送端提供了三种消息投递方式,分别为:同步发送、异步发送和单向发送。
(5.1.1) 同步发送
同步发送是RocketMQ默认的发送⽅式,默认调⽤MQProducer类的send⽅法,每次发送都需要等待响应,配合重试机制这种发送⽅式的可靠性是极⾼的,缺点当然是吞吐量低了。同步发送消息默认失败重试次数为2次,这个参数由DefaultMQProducer类中的参数retryTimesWhenSendFailed来控制。

当然我们可以在初始化生产者时,自行去配置修改这个参数:
DefaultMQProducer producer = new DefaultMQProducer();
producer.setNamesrvAddr("127.0.0.1:9876");
producer.setInstanceName("INSTANCE");
producer.setProducerGroup("PRODUCER-GROUP");
producer.setVipChannelEnabled(false);
producer.setMaxMessageSize(4194304);//单次消息发送内容小大
producer.setRetryTimesWhenSendFailed(3);//同步发送时消息失败重试次数
producer.setSendMsgTimeout(5000);//消息发送超时时间
(5.1.2) 异步发送
异步发送指的是⽣产者发送消息完后,不⽤同步等待RocketMQ服务端返回ACK确认标识,而是通过回调函数处理响应结果,因此需要我们自定义回调函数以及处理流程,以下是异步发送演示代码:
DefaultMQProducer producer = new DefaultMQProducer();
producer.setNamesrvAddr("127.0.0.1:9876");
producer.setInstanceName("INSTANCE");
producer.setProducerGroup("PRODUCER-GROUP");
producer.setVipChannelEnabled(false);
producer.setMaxMessageSize(4194304);//单次消息发送内容小大
producer.setRetryTimesWhenSendAsyncFailed(2);//异步发送时消息失败重试次数
producer.setSendMsgTimeout(5000);//消息发送超时时间
Message message = new Message();
message.setTopic("订阅的主题");
message.setTags("TagA");//设置消息Tag
message.setBody("消息体".getBytes());
producer.send(message, new SendCallback() {
@Override
public void onSuccess(SendResult sendResult) {
System.out.printf("消息【%n】发送成功", sendResult.getMsgId());
}
@Override
public void onException(Throwable e) {
System.out.printf("消息发送失败,原因 %n", e);
}
});
异步发送默认重试次数为2次,可以通过配置setRetryTimesWhenSendFailed来配置重试次数。
(5.1.3) 单向发送
单向发送是指不关⼼发送结果的发送⽅式,因其不关⼼发送结果故⽽也不⽀持重试,很少有业务会使⽤单向发送。
DefaultMQProducer producer = new DefaultMQProducer();
producer.setNamesrvAddr("127.0.0.1:9876");
producer.setInstanceName("INSTANCE");
producer.setProducerGroup("PRODUCER-GROUP");
producer.setVipChannelEnabled(false);
producer.setMaxMessageSize(4194304);//单次消息发送内容小大
producer.setSendMsgTimeout(5000);//消息发送超时时间
Message message = new Message();
message.setTopic("订阅的主题");
message.setTags("TagA");//设置消息Tag
message.setBody("消息体".getBytes());
producer.sendOneway(message);//单向发送
(5.2) 消费端的消息重试
我们刚刚了解了生产端的消息重试,现在如果是消费端的消费者出现了异常,消费某条消息失败后,RocketMQ应该如何进行消费重试呢?当消费者在消费消息失败时,RocketMQ会自动进行消息重试,默认最多会重复重试16次。如果一个条信息在重复重试16次以后还是无法被正常消费,这种情况极大概率是业务上的数据问题。此时RocketMQ会将消息直接扔到死信队列种,那么死信队列是什么呢?
死信队列⽤于处理⽆法被正常消费的消息。当⼀条消息初次消费失败,消息队列会⾃动进⾏消息重试;达到最⼤重试次数后,若消费依然失败,则表明消费者在正常情况下⽆法正确地消费该消息,此时,消息队列不会⽴刻将消息丢弃,⽽是将其发送到该消费者对应的特殊队列中。
RocketMQ 将这种正常情况下⽆法被消费的消息称为死信消息(Dead-Letter Message),将存储死信消息的特殊队列称为死信队列(Dead-Letter Queue)。在 RocketMQ 中,可以通过使⽤ console 控制台对死信队列中的消息进⾏重发来使得消费者实例再次进⾏消费。
(5.2.1) 消费重试策略
RocketMQ 的消费重试策略的内部机制其实是通过控制⼀下3个⼤将展开的:
- 重试过程状态机:控制消息在重试流程中的状态和变化逻辑
- 重试间隔:上⼀次消费失败或超时后,下次重新尝试消费的间隔时间
- 最⼤重试次数:消息可被重试消费的最⼤次数
通常是通过控制重试次数和重试时间间隔,来修改具体的重试状态来达到重试的⽬的。RocketMQ会为每个消费组都设置⼀个Topic名称为“%RETRY%+consumerGroup”的重试队列,重试队列是以消费者来进行创建,而不是根据主题创建。重试队列⽤于暂时保存因为各种异常⽽导致Consumer端⽆法消费的消息,每个Consumer实例在启动的时候就默认订阅了该消费组的重试队列Topic。考虑到异常恢复起来需要⼀些时间,会为重试队列设置多个重试级别,每个重试级别都有与之对应的重新投递延时,重试次数越多投递延时就越⼤。重试消息的处理是先保存⾄Topic名称为SCHEDULE_TOPIC_XXX的延迟队列中,后台定时任务按照对应的时间进⾏Delay后重新保存到相应地重试队列中。
- 最⼤重试次数
RocketMQ默认允许每条消息最多重试16次,可以通过客户端参数DefaultMQPushConsumer.maxReconsumeTimes设置最⼤重试次数,超过最⼤重试次数还消费失败,消息将会被⽆情的丢弃到死信队列中了,maxReconsumeTimes默认值为-1,-1代表重试16次。 - 重试间隔时间
当然,还可以通过控制重试间隔时间来⾃定义重试策略,重试间隔时间指的是,过多久重试⼀次。对于顺序消息,重试间隔为固定时间,默认 3000 毫秒;而对于⽆序消息,重试时间间隔是阶梯时间:

这样做是因为,通常故障恢复需要⼀定的时间,如果不间断的重试,重试⼜失败的情况会占⽤并浪费资源。
六、RocketMQ之消息如何存储以及处理
现在我们已经了解了RocketMQ的Topic主题、消息模型、消息过滤以及消息重试等基础知识,我们将继续探讨消息存储相关的知识点。消息队列一大核心功能是“削峰填谷”。在分布式系统中,下游(接收方)的处理能力通常是有上限的,在大促、秒杀或突发新闻时,瞬间流量可能是平时的几十倍。如果请求直接全部涌入后端,会瞬间压垮数据库,导致系统雪崩。这个突然的流量爆发点,我们通常称为“峰值(Peak)”,而在深夜或非高峰期,系统的处理能力是过剩的,资源处于闲置状态,我们称这个点为“谷值(Valley)”。
削峰(Smoothing the Peak):
当流量激增时,上游系统(如 Web 界面)不再直接请求数据库,而是将请求封装成消息丢进MQ。
逻辑: MQ 就像是一个巨大的容器,它能以极高的性能承接瞬间的冲击流量。对于用户来说,请求已经提交成功了(由于 MQ 写入极快,响应时间缩短)。
结果: 成功保护了后端脆弱的资源,避免了因过载导致的崩溃。
填谷(Filling the Valley):
后端服务(消费者)根据自己预设的节奏(如每秒处理 1000 条),匀速地从 MQ中拉取消息进行处理。
逻辑: 当高峰流量过去后,MQ 中依然积压着大量未处理的消息。后端服务会利用“流量低谷”的时间,继续消化这些积压的消息。
结果: 将原本集中的压力分散到了更长的时间维度上,使得系统资源得到了最充分的利用。
通常我们为了防止数据丢失,很多的做法都是将数据直接存放到数据库中,比如MySQL、Oracle、Redis等。但对于消息队列来说,其目的是为了削峰填谷,能够应对百万级的流量,而数据库在这方面的性能远远达不到。那么消息到底存在哪里呢?消息队列的消息是直接存储在磁盘上的,并不依赖第三方系统工具,不同的消息中间件,都按照自己定义的消息文件格式将消息直接存储到磁盘上。

(6.1) commitlog与consumequeue
你可能会问消息存储到磁盘上的话,那应该如何存储呢?假如有很多个 topic,每个 topic ⼜有很多个队列,是每个 topic ⼀个⽂件?还是每个队列⼀个⽂件? 其实为了保证RocketMQ 的写性能,最好是将所有消息都放在⼀个⽂件,并顺序读、顺序写,如果不放在⼀个⽂件,硬盘的存储物理位置就不是连续的,能⽆法保证顺序写,这个时候写⼊后的物理位置其实就是随机的了。
commitlog 是 RocketMQ 中唯一存储消息原始内容的文件,不管消息属于哪个 Topic、哪个队列,都会按「写入时间顺序」追加到 commitlog 文件中(顺序写,性能极高)。这个文件是以二进制来进行消息存储,使用二进制存储的好处在于可以节省存储空间,也能提高读取性能,其具有以下特性:
- 存储内容:消息的完整内容(包括主题、标签、消息体、属性、偏移量等)
- 文件格式:默认每个 commitlog 文件大小为 1GB,文件名为起始偏移量(如 00000000000000000000、00000000001073741824),写满一个就新建下一个
- 核心优势:顺序写磁盘(磁盘顺序写性能远高于随机写),这是 RocketMQ 高吞吐量的核心原因之一
- 存储位置:默认在 RocketMQ 安装目录的 store/commitlog下

细心的你可能还发现上图中有一个名叫consumequeue的文件夹,那这⾥⾯存的是什么呢?
consumequeue 翻译为「消费队列」,是 RocketMQ 为每个「Topic + 队列 ID」创建的索引文件,它不存储消息本身,只存储消息在 commitlog 中的「定位信息」,目的是让消费者能快速找到目标消息。
存储内容:每条记录仅包含 20 字节,由三部分组成:
8字节:消息在 commitlog 中的起始偏移量(定位消息物理位置);
4字节:消息长度(读取多少字节);
8字节:消息的 Tag 哈希值(快速过滤指定 Tag 的消息);
存储位置:默认在 store/consumequeue/{Topic}/{队列ID} 下(比如store/consumequeue/order_topic/0)
二者协同工作流程大致如下:
生产者发送消息:
1.消息被写入 commitlog(顺序写);
2.RocketMQ 的「刷盘线程」同步生成这条消息的索引记录,写入对应 Topic + 队列的 consumequeue 中;
消费者消费消息:
1.消费者根据自己订阅的 Topic 和分配的队列,读取对应的 consumequeue 索引;
2.根据索引中的「commitlog 偏移量 + 消息长度」,到 commitlog 中精准读取消息内容;
3.消费完成后,记录 consumequeue 的消费偏移量(下次从该位置继续消费)。
(6.2) RocketMQ之如何实现顺序消息
顺序消息是业务中常⽤的功能之⼀,⽽ RocketMQ 默认发送的事普通⽆序的消息,要保证消息的顺序,要从⽣产端到 broker 消息存储,再到消费消息都要保证链路的顺序,才可以做到真正的顺序消息。
顺序消费从严格划分来说分为普通顺序和严格顺序。普通顺序是指的相同队列收到的消息是有序的,通常也称为局部有序,条件是消息必须属于同一个队列,而队列本身就是先进先出的原则,因此自带顺序。而严格顺序就是不管消息属于哪个队列或者哪个主题,都必须满足先进先出的原则。

要保证顺序消息,首先要保证生产者发送的消息是顺序的!要保证发送的消息的顺序,只能采用同步发送并且是单个生产者进行消息发送,这是两个必要条件!究其原因还是因为顺序发送的技术原理,其技术原理也⽐较简单,就是要将同⼀类消息发送到相同队列即可,同一类消息在发送到MQ之前已经确保了发送顺序(同步、单个生产者发送),当消息存储到队列后,由于队列本身就是先进先出的原则,因此就完成了发送有序。但是如何将同一类的消息指定发送到同一个队列呢?RocketMQ提供了MessageQueueSelector接口,允许开发者自定义队列选举策略,默认的SelectMessageQueueByHash是根据参数的hashCode与当前主题队列数取模得到一个直,然后拿这个值作为队列数组下标获取队列。
public class SelectMessageQueueByHash implements MessageQueueSelector {
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
int value = arg.hashCode();
if (value < 0) {
value = Math.abs(value);
}
value = value % mqs.size();
return mqs.get(value);
}
}
在发送消息的时候,⽣产者可以参考这个实现逻辑,自定义选择器来选择指定的队列:
sendResult = producer.send(msg, new MessageQueueSelector() {
@Override
public MessageQueue select(List<MessageQueue> mqs, Message msg, Object arg) {
Integer hashCode = (Integer) arg;
int index = hashCode % mqs.size();
return mqs.get(index);
}
}, bizCode.hashCode());
- bizCode指的某种业务类型的编码,取bizCode的hash值作为选择器的参数
- ⽤哈希值和队列数 mqs.size() 取摸,得到⼀个索引值,结果会⼩于队列数
- 根据索引值从队列列表中取出⼀个队列,那 hash 值相同,取出的队列也就相同
通过发送消息的时候指定队列⽅式就保证了同⼀个业务类型消息发往相同的队列,也即保证消息发送的顺序性。
在确保了顺序发送和顺序存储后,现在要解决的就是顺序消费。RocketMQ支持两种消费模式,集群消费和广播消费。集群消费模式下每一条消息只会被ConsumerGroup分组下的一个Consumer消费,而广播消费模式下,每个Consumer都会消费这条消息。多数场景下⽤的都是集群消费,也就是⼀次消费代表⼀次业务处理,每⼀条消息都将由集群中
的⼀个实例来对应处理,顺序消费也叫有序消费,如果消息是顺序发送,且顺序存储,那理应消费也是按顺序⼀条条消费,但实际却没这么简单。
在ConsumerGroup中不⽌⼀个线程在那消费,因为同⼀个消费者可能会处理不同的队列消息,如果只有⼀个线程。那不得慢死,实际上会有多个线程同时消费,对应的是ConsumerGroup中的消费线程池。当涉及多线程进行操作时,如何防止消息被多个线程重复消费是一个棘手点。在RocketMQ中使用了3把锁来防止消息被重复消息,分别是分布式锁、Synchronized、ReentrantLock。
更多推荐
所有评论(0)