MQ进阶 - 分布式下的挑战
上文讲了三大主流 MQ 的的功能和特性,那就像是工具箱,在实际场景中按需使用即可。
本文讲述的是使用 MQ 我们会遭遇的问题,相当于 BUG —— 使用 MQ 出现了不符合预期的结果,不解决不可上线。
引言 - 分布式的经典问题
在分布式场景下,给我们带来了高可用、易扩展和高容错的便利,也带来了分布式系统下不少的麻烦,增加了分布式开发的复杂性。
分布式的问题一图概览:

可知主要问题是三个根源引起:网络不可靠、节点独立性和无全局时钟。
1)网络不可靠:在我们物理世界,网络波动(延迟、丢包、分区甚至中断)是无法根除的情况,那么这会导致:
- 消息丢失或重复
- 网络分区: 分布式系统中,不同节点之间因为网络故障,无法互相通信,从而被“分隔成多个孤岛”。 某个节点故障必然形成网络分区。
2)节点独立性:分布式系统由多个独立计算机(节点)组成,每个节点都有自己的内存、CPU和时钟。独立节点可能会出现:
- 节点故障:任何节点都可能随时发生故障(宕机、网络中断、磁盘损坏等),这与单机程序中“要么全好、要么全坏”的模式截然不同,系统必须能处理部分失效(Partial Failure)的情况。
3)无全局时钟:虽然每个节点都有自己的物理时钟(如石英钟),但它们之间无法做到精确同步。
- 时钟不同步:即使通过NTP(网络时间协议)同步,也存在毫秒级甚至更大的误差,且时钟会因温度、电压等因素漂移
- 事件顺序难定:因为时钟不同步,节点A记录的事件时间戳可能早于节点B,但实际发生顺序却相反。这使得确定跨节点事件的绝对顺序变得不可能
为什么大家使用北京时间,还会出现时间不同步呢?
服务器通常会定期通过 NTP 同步时间到“标准北京时间”(如国家授时中心)。 但是同步不是实时的,比如每隔几分钟或几小时才同步一次。 而网络延迟、抖动、负载都会导致同步结果存在 误差(几十毫秒到几百毫秒)。
另外 每台机器内部都有一个石英振荡器来计时。 操作系统每隔一段时间会读这个振荡器的值来更新系统时间。 但是!不同机器的振荡器频率略有差异(例如每秒快几微秒、慢几微秒),时间会逐渐漂移(drift)。
以上都是在分布式系统下无法突破的物理层面限制,这些物理层面的问题,我们必须在软件层面克服,物理层面无法突破,因此下列问题是软件开发者必须面对的经典挑战:
| 问题类型 | 具体表现 | 根源关联 |
|---|---|---|
| 数据一致性 | 数据在多个副本间可能出现不一致。 | 网络分区、节点故障、消息延迟与重复。 |
| 可用性 | 某节点失联或网络中断后,服务完全中断或响应严重变慢 | 节点独立性、网络不可靠 |
| 分布式共识 | 多个节点如何对某个值(如谁是Leader)达成一致 。 | 节点独立性、网络不可靠。 |
| 分布式事务 | 确保跨多个节点的操作序列要么全部成功,要么全部失败 ,分布式难以实现 | 网络不可靠 、节点故障、无全局时钟。 |
| 故障检测与恢复 | 如何准确判断一个节点是真的宕机了还是只是网络慢,以及如何恢复。 | 网络不可靠(延迟与分区)、节点独立性。 |
| 消息可靠性(丢失/重复) | 消息在网络中丢失、延迟或被重复投递,导致消费者状态错误 | 网络不可靠、节点独立性 |
将分布式的挑战讲的深入一点,是为了看清楚消息队列的重复消费、幂等性、一致性问题是由分布式架构带来的,根源是出自物理层面不可避免的问题,我们来学习前辈如何巧妙的在软件层面攻破物理层面的困难。
| 问题类型 | 具体表现 | 常见解决方案 | RabbitMQ 支持情况 | RocketMQ 支持情况 | Kafka 支持情况 |
|---|---|---|---|---|---|
| 消息丢失 | 生产者发送失败、Broker 存储失败、消费者处理失败但未确认,导致消息永久丢失。 | 生产者确认机制、Broker 持久化、消费者手动确认、副本机制 | ✅(发布者确认、持久化队列/消息、手动确认)rabbitmq.comrabbitmq.com | ✅(同步/异步发送、刷盘策略、主从复制)AlibabaCloudStackademic | ✅(acks = –1(所有副本确认)、副本机制、消费者位移提交) |
| 重复消费 | 网络延迟或确认失败导致消息被重投,或消费者重启后重新消费,导致同一消息被处理多次。 | 消费端幂等设计、唯一ID去重(DB/Redis) | ❌(主要依赖业务设计) | ❌(推荐业务唯一标识) | ✅(生产者幂等(单分区/单会话)、流处理 Exactly Once 支持;但消费端仍需幂等处理) |
| 消息顺序 | 需要保证一组有序的消息(如同一订单的创建→支付→发货)按照发送顺序被消费。 | 全局有序(单队列)或分区有序(按 key 路由同一队列) | ✅(队列内 FIFO;全局有序需特殊配置) | ✅(分区有序:顺序消息支持相同 Key 到同一队列) | ✅(分区内有序:相同 Key 发到同一分区) |
| 消息积压 | 生产者速度远超消费者,或消费者故障,导致消息在 Broker 中大量堆积。 | 水平扩展(增加消费者/分区)、监控告警、限流降级 | ✅(监控队列长度、可增加消费者) | ✅(高性能堆积场景支持、监控与扩展能力强) | ✅(天然支持海量堆积、可增加分区和消费者) |
| 事务消息 | 跨多个系统/数据库的操作必须具有原子性(例如订单创建成功才发送消息通知下游)。 | 两阶段提交(2PC)、本地消息表、事务消息机制 | ❌(原生不支持,仅可通过插件或本地消息表+Confirm 模拟) | ✅(原生支持事务消息机制:半消息 → 提交/回滚) | ✅(支持流处理场景的事务/Exactly Once,但并非传统队列事务机制) |
| 分布式事务 | 涉及多个消息生产与消费的完整业务流程需保证最终一致性。 | Saga 模式、TCC 模式、可靠消息最终一致性 | ❌(主要依赖业务模式:Saga 或本地消息表结合) | ✅(提供基础机制:事务消息为可靠消息最终一致性利器) | ✅(提供基础机制:在流处理里支持事务(EOS)) |
| 重试补偿 | 消费失败后自动/手动重试,直到达到上限或人工干预(补偿)处理失败消息。 | 自动重试机制、死信队列、人工补偿 | ✅(原生支持死信队列,消息被拒绝或过期时可路由至DLX) | ✅(原生支持消费重试,达到上限后进入死信队列) | ❌(无内置重试机制,需业务方自行实现,如发到重试Topic) |
- 事务消息:主要针对的是“本地事务操作与消息发布之间的一致性问题”,也就是当你做操作(例如数据库写入)然后发送消息通知下游时,如何做到“要么操作+消息都成功,要么都失败”。
- 分布式消息:泛指在分布式系统中各节点之间通过消息通信所涉及的“异步、跨节点、解耦、可靠性”问题。它更宽泛,不仅包括事务相关,还包括消息传递、失败重试、顺序、分区、幂等等。
在进入这些问题的学习之前,我们再了解一个分布式下的必要知识点,CAP定理。
CAP 定理
CAP : 一致性(Consistency)、可用性(Availability)、分区容错性(Partition tolerance)
根据百度百科:
CAP定理的思想由加州大学伯克利分校的计算机科学家埃里克·布鲁尔 于2000年在分布式计算原则研讨会(PODC)上首次提出,当时被称为CAP猜想。
2002年,麻省理工学院(MIT)的赛斯·吉尔伯特和南希·林奇 发表了正式证明,从理论上验证了这一猜想的正确性,使其成为一个定理。
CAP是什么?
在分布式系统健康的状态下,不存在CAP理论,因为分区容错性,一致性和可用性完全可以得到保障。当单个节点发生故障时,必须优先保证分区容错性,此时 CAP 定理发挥作用。
CAP定理指出,任何一个分布式系统(由多个通过网络通信的节点组成的系统)无法同时同时满足以下三个基本需求,最多只能同时满足其中的两个
- 一致性: 所有节点在同一时间看到的数据是一致的。
- 可用性: 每个请求都会得到响应(不管成功或失败)。
- 分区容错性: 系统能继续运行,即使网络部分节点无法通信。
CAP 定理揭示了分布式系统中的一种固有状况和必然的权衡,一种无法解决的困境。
为什么不能三个都满足
理想情况下,我们肯定希望三个都满足,如果不能,我们必须要优先满足分区容错性,因为网络故障是不可避免的,如果分区容错性不能保证,那么其它两个更加做不到。
如果选择保证一致性(C):那么在分区期间,某些节点可能需要停止服务(拒绝读写请求),以避免数据不一致,这就牺牲了可用性(A)
如果选择保证可用性(A):那么系统即使在分区期间也必须响应每个请求,但这可能意味着返回的是旧数据(不同节点数据可能不一致),这就牺牲了一致性(C)
举例论证
系统故障,发生以下网络分区的情况:
- 假设系统有三个节点 A、B、C
- 节点 A 和 B 可以互相通信。
- 节点 C 与 A/B 网络断开(分区)。
- 现在系统出现了孤岛:A/B 和 C 之间无法通信。
如果选择一致性,可用性必然下降:
- 系统必须保证每次写入都被所有节点确认(或者按某种算法保证一致性)。
- 但是节点 C 无法与 A/B 通信。
- 如果继续让 C 提供写入或读取,就可能看到不同步的数据 → 破坏一致性。
- 所以为了 保证一致性,节点 C 就必须 拒绝读写请求,相当于停服。
如果**选择可用性,**一致性必然下降:
- 即使 C 与 A/B 失去连接,C 仍然响应请求(读/写)。
- 但是 C 的数据可能与 A/B 不一致(写入只在 C,未同步到 A/B)。
- 这就破坏了全局一致性(C)。
强行让 C 服务器运转,造成数据不一致。
消息队列的基本流程
我们先看看 RocketMQ 和 RabbitMQ 生产消息和消费消息的最基本流程。
RocketMQ
核心角色
- Producer (生产者): 消息的发送方。负责产生消息,并发送到指定的 Topic。
- Consumer (消费者): 消息的接收方。负责从指定的 Topic 订阅并消费消息。
- Broker (代理服务器): 消息中转角色。负责存储来自生产者的消息,并将消息投递给消费者。是消息存储和转发的核心。
- NameServer (注册中心): 无状态的轻量级服务。充当服务发现的角色,维护了 Broker、Topic 等路由信息。生产者和消费者会从 NameServer 获取目标 Broker 的地址列表。
RocketMQ 生产消息到 Brocker 的详细流程:
- 创建生产者: 应用程序创建一个
**DefaultMQProducer**实例,并设置 NameServer 的地址。 - 启动生产者: 调用
**producer.start()**方法。
-
- 内部会启动多个服务线程。
- 与 NameServer 建立长连接,定期(默认每30秒)从 NameServer 拉取最新的 Topic 路由信息(即该 Topic 对应的 Broker 地址、队列等)。
- 构建消息: 创建一个
**Message**对象,包含 Topic、Body(消息体)和可选的 Tag/Keys。 - 发送消息: 调用
**producer.send(message)**方法。
-
- 选择队列: 生产者客户端会根据预设 (代码逻辑预设) 的负载均衡策略(如轮询),从获取到的路由信息中选择一个消息队列(Message Queue,位于某个 Broker 上)。
- 网络发送: 将消息发送给目标 Broker。
- 等待 Broker ACK (关键步骤):
-
- 默认情况下,
**send()**方法是同步阻塞的。它会一直等待,直到收到来自 Broker 的响应,或者直到超时。 - Broker 接收到消息后,会将其写入 CommitLog(内存和磁盘),然后向生产者返回一个
**SendResult**对象。 - 这个
**SendResult**就是生产端的 ACK。它包含了发送状态(**SendStatus**),如**SEND_OK**(成功)、**FLUSH_DISK_TIMEOUT**(刷盘超时)、**FLUSH_SLAVE_TIMEOUT**(同步到从机超时)、**SLAVE_NOT_AVAILABLE**(从机不可用)等。
- 默认情况下,
- 处理结果: 生产者根据
**SendResult**的状态,可以判断消息是否成功发送,并进行相应的业务处理(如重试、记录日志等)。

RocketMQ 生产消息到 Brocker 的详细流程:
- 创建消费者: 应用程序创建一个
**DefaultMQPushConsumer**实例,设置 NameServer 地址、消费者组名(Consumer Group)。 - 订阅 Topic: 调用
**consumer.subscribe("TopicName", "\*")**方法,订阅感兴趣的 Topic。 - 注册消息监听器: 调用
**consumer.registerMessageListener()**注册一个监听器实现(例如**MessageListenerConcurrently**)。这个监听器包含了业务处理逻辑。 - 启动消费者: 调用
**consumer.start()**方法。
-
- 内部会与 NameServer 建立长连接,获取 Topic 路由信息和 Broker 地址。
- 根据消费者组名和负载均衡策略,为当前消费者分配若干个消息队列进行消费。
- 启动后台拉取线程,主动向其负责的 Broker 拉取消息。
- Broker 投递消息: Broker 的拉取处理器接收到消费者的请求后,从存储中读取消息并返回给消费者,如果没有消息会进行长轮询。
- 消费消息 (业务处理):
-
- 消费者客户端收到消息后,会调用用户注册的
**MessageListener**的**consumeMessage**方法,将消息交给业务代码处理。
- 消费者客户端收到消息后,会调用用户注册的
- 返回消费状态 (关键步骤):
-
- 业务代码处理完毕后,
**consumeMessage**方法必须返回一个**ConsumeConcurrentlyStatus**状态。 **CONSUME_SUCCESS**: 表示消费成功。**RECONSUME_LATER**: 表示消费失败,希望稍后重新消费。
- 业务代码处理完毕后,
- 发送消费端 ACK (关键步骤):
-
- 如果
**consumeMessage**返回**CONSUME_SUCCESS**:消费者客户端会自动向 Broker 发送一个 ACK 请求。 - 如果
**consumeMessage**返回**RECONSUME_LATER**或抛出异常:消费者客户端不会发送 ACK。
- 如果
- Broker 处理 ACK:
-
- Broker 收到 ACK 后,会将该消息在消费队列中的偏移量更新,标记为已消费。
- 如果在一定时间内(默认约15分钟)未收到某条消息的 ACK,Broker 会认为消费失败,并重新将该消息投递给消费者组内的其他消费者(如果有的话),或在下次拉取时再次投递给原消费者。

RabbitMQ
核心角色
- Producer (生产者): 消息的发送方。负责将消息发送到交换机。
- Consumer (消费者): 消息的接收方。负责从队列中获取消息并进行处理。
- Broker (代理服务器): RabbitMQ 服务本身。负责接收、存储和路由消息。
- Exchange (交换机): 接收来自生产者的消息,并根据路由键将消息推送到一个或多个队列。生产者从不直接将消息发送到队列。
- Queue (队列): 存储消息的缓冲区,直到消费者将其取走。
- Channel (信道): 位于 TCP 连接内部的虚拟连接。它是进行 AMQP 操作(如发布、消费、获取消息)的轻量级载体,避免了为每个操作都建立 TCP 连接的开销。
RabbitMQ 生产消息到队列的详细流程:
- 创建连接和信道: 应用程序与 Broker 建立一个 TCP 连接,并在其上创建一个信道。
- 启用发布者确认 (关键步骤): 在信道上调用
**channel.confirmSelect()**,将该信道设置为确认模式。这是开启生产端 ACK 的前提。 - 声明交换机和队列: 为了保证消息能被正确路由,生产者通常会声明一个交换机和一个队列,并将它们绑定起来(这是一个幂等操作,重复声明不会出错)。
- 发送消息: 调用
**channel.basicPublish()**方法,将消息发送到指定的交换机,并附带路由键。 - 等待 Broker ACK (关键步骤):
-
- 消息发送后,生产者不会立即收到响应。Broker 接收到消息后,会将其路由到匹配的队列。
- 一旦消息至少被一个队列接收并持久化(如果队列是持久化的),Broker 就会通过该信道向生产者发送一个确认帧,包含消息的
**delivery-tag**。 - 如果由于某些原因(例如找不到匹配的队列)消息无法被路由,Broker 会发送一个未确认帧。
- 处理确认: 生产者可以监听信道的确认和未确认事件,根据收到的
**delivery-tag**来判断哪条消息成功或失败,并进行相应的补偿业务处理(如记录日志、重发失败的消息等)。

RabbitMQ 消费流程:
- 创建连接和信道: 与生产者类似,消费者也需要建立连接和信道。
- 声明队列: 消费者声明它要消费的队列,确保队列存在。
- 设置 QoS (可选但推荐): 调用
**channel.basicQos(prefetchCount)**,限制消费者每次能从队列预取的消息数量。这可以防止消费者被消息淹没,实现更公平的负载均衡。 - 定义消费者并消费: 创建一个
**Consumer**实现类(或使用回调),并调用**channel.basicConsume()**方法开始消费。
-
- 关键参数: 在
**basicConsume**中,必须设置**autoAck**为**false**。这表示手动确认模式。如果设置为**true**(自动确认),RabbitMQ 会在消息发送给消费者后立即删除它,不管消费者是否处理成功,这是极不可靠的。
- 关键参数: 在
- Broker 投递消息: RabbitMQ 将队列中的消息推送给消费者。
- 消费消息 (业务处理): 消费者的回调方法被触发,消息内容被传递给业务代码进行处理。
- 发送消费端 ACK (关键步骤):
-
- 如果业务处理成功: 在
**handleDelivery**回调方法的最后,调用**channel.basicAck(deliveryTag, multiple)**。**deliveryTag**是消息的唯一标识,**multiple**表示是否确认小于该**deliveryTag**的所有消息。 - 如果业务处理失败: 调用
**channel.basicNack(deliveryTag, multiple, requeue)**或**channel.basicReject(deliveryTag, requeue)**。**requeue**参数决定消息是否重新返回队列头部等待再次投递。
- 如果业务处理成功: 在
- Broker 处理 ACK:
-
- Broker 收到
**basicAck**后,才会从队列中移除该消息, 会删除持久的消息。 - 如果收到
**basicNack**/**basicReject**且**requeue=true**,消息会被重新放回队列。 - 如果消费者在处理消息时断开连接且未发送 ACK,Broker 也会认为消费失败,并将消息重新入队。
- Broker 收到

消息丢失
消息丢失指的是:生产者发送了消息,但这条消息未能被成功存储或被消费者处理,最终在系统中“消失”,不再被传递或消费。消息丢失问题会严重影响数据一致性和业务可靠性
问题产生
消息丢失通常有以下几个环节可能发生:
- 生产者发送失败:生产者向 Broker 发送消息时,可能因为网络中断、Broker 不可用、客户端超时而失败。如果生产者误以为发送成功但实际上没有,则消息丢失。
- Broker 存储失败或确认失败:消息到达 Broker,但在持久化(写磁盘)或写入日志、刷盘、同步副本等过程中失败,或者 Broker 挂掉,还未完成持久化。
- 消费者处理失败并且未重试/未记录:消费者接收消息后处理失败,且系统没有做重试或者死信机制、也没有记录该消息状态,导致消息被“吞掉”。
- 配置或使用不当:例如队列或主题未配置持久化、生产者未启用确认机制、Broker 副本机制配置不当、网络故障未被捕捉等。
总结一句话:消息丢失的根因通常是网络、存储、确认机制不可靠,原因很多,无法全部概括。
问题解决
确保消息不丢失需要从生产者、Broker、消费者三个层面进行保障
- 生产者端:
-
- 使用发送确认机制:确保消息成功发送到Broker。
- 配置重试机制:发送失败时自动重试
- Broker端:
-
- 消息持久化:将消息和队列设置为持久化,写入磁盘
- 配置副本机制:通过主从复制或多副本保证数据冗余
- 消费者端:
-
- 手动提交offset:在消息处理完成后手动确认(程序员代码中用业务逻辑控制,不依赖消息中间件提供的自动提交),避免自动提交导致的问题
- 消费幂等性设计:确保重复消费不会导致业务逻辑错误
但是这不能百分百保证消息丢失,所以下一个小节会讲消息补充,我们能接受极端情况下的丢失。
RabbitMQ 解决
RabbitMQ通过多层次确认机制和持久化策略保障,其核心思想是在消息传递的每个关键节点设置确认点。
- 生产者确认机制
-
- Publisher Confirms: 当生产者发送消息到RabbitMQ后,Broker会返回一个确认结果(ACK / NACK)
- Publisher Returns : RabbitMQ 提供了 ReturnCallback 回调机制 ,开启mandatory=true 发送消息 时, 如果交换机找不到任何队列匹配,Broker 会 立即回调生产者,告诉你消息无法路由。 这样生产者就知道有消息没被队列接收,可以选择 重试、记录日志或发送到死信队列。
- 持久化策略
为了保证消息不丢失,RabbitMQ不仅仅持久化了消息,还持久化了 Exchange 和 Queue , 做到了三层立体式持久化。
| 层级 | 持久化对象 | 含义 | 作用 |
|---|---|---|---|
| ① Exchange 层 | 交换机持久化 | Broker 重启后,交换机配置(名称、类型、绑定关系)仍然存在 | 确保消息路由基础结构不丢 |
| ② Queue 层 | 队列持久化 | Broker 重启后,队列仍然存在(包括绑定信息) | 确保消息能有地方落地 |
| ③ Message 层 | 消息持久化 | 消息内容写入磁盘(而不是只在内存) | 确保消息数据不丢失 |
- 消费者确认机制
消息确认是消费者向Broker发送确认,表示这条消息处理完了, Broker 收到 ack 后,才会把这条消息从队列里删除。
**RabbitMQ 支持自动确认,**缺点是容易丢失消息,因为自动确认机制是在消费者真实消费数据之前发送 ACK 。自动确认吞吐量高。
手动确认:可以自己写代码,灵活控制发送ACK的时机,更稳定。
RocketMQ 解决
- 生产者可靠性机制
- 同步发送:使用同步发送方式,等待Broker返回确认结果,例如
**SendResult sendResult = producer.send(msg)**。通过检查**sendResult.getSendStatus()**判断是否发送成功,失败时可进行重试 - 事务消息:对一致性要求极高的场景,使用事务消息。它通过两阶段提交(2PC)确保本地事务与消息发送的原子性
- Broker端可靠性机制
- 同步刷盘(SYNC_FLUSH):消息写入内存后立即刷盘,确保数据落盘,防止因Broker突然宕机导致内存中未刷盘的消息丢失
- 异步刷盘(ASYNC_FLUSH):消息写入内存后异步刷盘,性能更好
- 主从复制:配置主从Broker。将
**brokerRole**设置为**SYNC_MASTER**(同步复制)。Master接收到消息后,会同步等待Slave复制成功,才向生产者返回ACK。即使Master宕机,数据也在Slave上有备份
- 消费者可靠性机制
- 消费成功后才提交Offset:消费者在业务逻辑执行成功后,再向Broker提交消费进度。如果处理过程中消费者宕机,消息会重新投递,不会丢失
- 避免异步消费:RocketMQ官方建议不要在消费端使用异步处理消息,以免因未处理完就提前提交Offset导致消息丢失

消息丢失问题不能根除(ack丢失,刷盘崩溃等), 接受极端情况丢失,但用补偿机制兜底
消息补偿
消息补偿是一种事后补救机制,用于在消息传递或处理失败时,通过重新发送或执行补偿操作,确保业务数据最终达到一致状态,它通常用于处理分布式事务中的异常情况,特别是在采用最终一致性模型的系统中。
具体来说,生产者发送消息失败、消费者消费失败、消息超时未确认、事务状态未知会触发消息补偿。
rabbitMQ 消息补偿机制
- 生产者确认机制
生产者发送消息时,可开启 Confirm 模式。Broker 会异步确认消息是否到达交换机,会出现两种情况:
- ACK:消息成功到达交换机。
- NACK:消息未到达交换机。生产者可在回调中捕获 NACK,执行重新发送或记录日志等补偿操作
- 死信队列
当消息成为“死信”(如消费者拒绝且不重新入队、消息过期、队列满)时,若配置了死信交换机(DLX),消息会被路由到死信队列。进入死信队列的消息,可在程序中写逻辑对死信队列进行重试、重发(让生成者再发一次)、记录日志、人工介入。
- 消息重试与消费者确认
消费者处理失败时,可拒绝消息并选择不重新入队(**requeue=false**),让消息进入死信队列,避免无限重试导致阻塞。同时,应确保业务逻辑的幂等性,以应对可能的消息重复
RocketMQ 消息补偿机制
RocketMQ 提供了更强大的事务消息机制,并依赖定时任务和状态回查来实现补偿。
- 事务消息
补偿机制(事务状态回查):若Broker长时间未收到确认,会反向回查生产者的事务状态,生产者需实现回查接口,根据本地事务状态决定提交或回滚。事务回查触发消息补偿。
- 定时任务与消息重试
生产者:发送失败时可自动重试,自带消息补偿。
消费者:消费失败时,消息会自动重试。超过最大重试次数后,消息会进入**死信队列,**可监控死信队列并进行补偿。
- 本地消息表
对于未使用事务消息的场景,可采用本地消息表方案
本地消息表: 在生成者端建立一个消息表,用来把消息存储在业务数据库中,和业务操作放在同一个本地事务里提交。 这样就能确保业务操作成功的同时,消息也一定成功发送,避免中间的缝隙发生其它意外导致业务操作成功,消息发送失败的情况。
本地消息表是一种补偿机制, 先落地保存“要发的消息”,等网络恢复、Broker可用,再补偿性地重发。 如果没有收到 Broker 的 ACK 确认,或者Broker 返回失败回调 ,本地消息表发挥作用。
重复消费
同一条消息被消费者处理了多次。
消息在 Broker(消息中间件)只存在一条,但消费者因为确认或网络问题,重复地拿到同一条消息去执行业务逻辑。
导致重复消费的原因:
- ACK 没发出去 ,Brocker 重试机制触发
- 重启后没来得及提交 offset 或确认状态,又从上次位置重新消费。
- 手动补偿/重放逻辑执行多次 , 运维或定时任务手动重发时,逻辑重复执行。
重复消费的危害
- 数据错误:比如重复创建订单、重复扣款,导致业务数据混乱
- 资源浪费:无谓地消耗CPU、内存和数据库连接等系统资源。
- 业务逻辑混乱:比如库存被多扣减,导致超卖
- 用户体验差:用户可能收到多条重复的短信或通知
解决重复消息的思路
无法避免消息不会重复发送,所以解决思路应该是消费的幂等性。
即无论消息被消费多少次,最终的结果都与消费一次相同。方案如下:
1. 唯一ID + 去重表
为每条消息生成一个全局唯一ID(如订单号、UUID),并在消费端创建一个去重表(通常使用数据库)。
- 消费前:消费者在消费前,先查询去重表中是否存在该ID。
- 如果存在:说明已处理,直接跳过。
- 如果不存在:执行业务逻辑,并将该ID插入去重表
▼sql复制代码-- 去重表示例 CREATE TABLE `message_consumed` ( `message_id` VARCHAR(64) PRIMARY KEY, `consume_time` TIMESTAMP DEFAULT CURRENT_TIMESTAMP );
2. Redis 原子操作
利用Redis的**SETNX**(SET if Not eXists)命令实现原子性去重。
- 消费前:尝试以消息ID为key向Redis写入一个值(并设置过期时间)。
- 写入成功:说明首次消费,执行业务逻辑。
- 写入失败:说明已消费,直接跳过
▼java复制代码// 伪代码示例 String messageId = message.getMessageId(); Boolean isSuccess = redisTemplate.opsForValue().setIfAbsent(messageId, "1", 30, TimeUnit.MINUTES); if (Boolean.TRUE.equals(isSuccess)) { // 首次消费,执行业务逻辑 processBusiness(message); } else { // 重复消息,直接跳过 logger.info("重复消息,已跳过: {}", messageId); }
3. 数据库唯一约束
对于业务数据本身有唯一标识的场景(如订单号),可以直接在数据库表中设置唯一索引或主键约束。
即使重复消费导致重复插入,数据库会因约束冲突而报错,应用捕获异常后忽略即可
▼sql复制代码-- 订单表order_id设置唯一约束 CREATE TABLE `orders` ( `order_id` VARCHAR(64) PRIMARY KEY, -- 其他字段... );
4. 业务状态机判断
对于有明确状态流转的业务(如订单状态:待支付→已支付→已发货),可以通过检查当前业务状态来判断是否需要处理。
如果消息要求的操作(如“支付”)对应的当前状态(如“已支付”)已经符合,则说明是重复或过时消息,直接丢弃
5. 消息添加辅助字段
在 RocketMQ 中,当生产者发送消息时,可以给每条消息指定一个“业务唯一键”,例如订单号、支付流水号:
▼java复制代码Message msg = new Message( "OrderTopic", "TagA", "KEY_1001", // 业务唯一键 "Order created".getBytes() ); producer.send(msg); // 消费端 message.getKeys()
消费者拿到 Key 后,可以在自己系统里判断该消息是否已处理过。
这个 key 是 RocketMQ 给你暴露的幂等辅助字段 ,通过它实现幂等。
RabbitMQ 消息结构更加通用,可以设置在 **message_id** 或 headers 里。
6. 乐观锁
在数据表中增加一个 **version** 字段(版本号)。更新数据时,需要带上这个版本号。
▼sql复制代码UPDATE products SET stock = stock - 1, version = version + 1 WHERE id = 123 AND version = 5;
应用在更新前先读取 **version=5**,执行更新时,在 **WHERE** 条件中带上这个版本号。然后检查受影响的行数:
- 如果行数为1,说明更新成功。
- 如果行数为0,说明在你读取和更新之间,已经有其他请求修改了这条数据(version已经不是5了),本次更新失败。
通过 version 字段来确保不会重复执行,更新数据成功 version + 1 , 这比数据库字段唯一约束更宽松,但实现复杂度高了一点。
消息顺序
消息顺序(Message Ordering)指的是消息按照生产者发送的先后顺序被消费者接收和处理。在消息队列中,这通常分为两种类型
- 全局有序:整个主题(Topic)中的所有消息,都严格按照生产者发送的先后顺序进行消费。例如,发送顺序是 M1, M2, M3,消费时也必须是 M1, M2, M3。
- 分区有序(消息键分区):只需要保证具有某个相同特征(如同一个订单ID、同一个用户ID)的消息能够被顺序消费即可。例如,订单
**A**的消息**A1, A2, A3**顺序消费,订单**B**的消息**B1, B2, B3**也顺序消费,但**A1**和**B1**谁先消费无所谓。
消息顺序的重要性
操作的顺序至关重要,乱序会导致业务逻辑错误或数据不一致。
比如订单的状态流转:创建订单 → 支付订单 → 发货 → 确认收货,一个状态是一条消息,上游执行成功会给下游发送一个消息,如果消息不会顺序消费,那么会发生数据不一致,业务混乱。
在这个场景中,每一个状态肯定是顺序发给队列的,只要确保同一个订单的消息进入同一个队列,那么在队列中的消息就是物理有序,RocketMQ 默认是消费组内单实例的单线程消费,不会造成乱序的问题。
作为开发者,我们只要保证同一个订单的消息发送给固定的一个队列就好,这就要用到消息键。
消息乱序的原因
消息队列(如 RabbitMQ、Kafka、RocketMQ、Redis Stream)都是并发消费模型。 如果你启用了多个消费者或多线程消费,一个消息队列中不同消息可能被不同线程同时处理。
实现消息有序
我们知道有消息顺序:全局有序和分区有序,其中分区有序是重点难点,使用的也是比较多的,因为分区有序性能高,但是复杂一点。
分区有序使用消息键模式实现,接下来看各个MQ怎么处理。
kafka 与 Rocket MQ 实现方案相同
消息键:在生成消息的时候指定 **一个能区分业务功能、能够分组业务的字段 ,**字段可以是任何类型,会经过哈希计算,对Brocker的总队列取模,相同的业务便投入到固定队列。
以kafka为语境进行讲解。
原理详解:
-
生产者将消息路由到固定分区(也就是队列)
-
- 定义业务键:为每条消息设定一个能够标识其业务分组的键。例如,处理订单消息,
**order_id**就是绝佳的业务键。 - 根据键计算分区:生产者在发送消息时,不是随机选择分区,而是根据业务键进行哈希计算,然后对分区总数取模,从而确定该消息应该发往哪个分区。
- 示例代码:
- 定义业务键:为每条消息设定一个能够标识其业务分组的键。例如,处理订单消息,
▼java复制代码// 业务消息 Message message = new Message("topic_order_paid", "订单A-已支付"); String businessKey = "order_A"; // 例如,订单ID // 计算目标分区 int partitionCount = 4; // 假设主题有4个分区 int targetPartition = Math.abs(businessKey.hashCode()) % partitionCount; // 发送到指定分区 producer.send(message, targetPartition);
效果:所有 **order_A** 的消息(如“已支付”、“已发货”、“已签收”)都会被哈希算法路由到同一个分区(比如分区 2)。同样,所有 **order_B** 的消息也会被路由到同一个分区(比如分区 0)。
-
消费者端单线程消费分区
-
- **天生支持单线程消费:**消息队列的消费者组模型天然支持这一点。一个分区在同一时间只能被消费者组里的一个消费者实例消费。
- 多个消费组可以订阅一个topic, 那么多个消息组都有自己的 offset , 一条消息会被多个消费组消费,这是跨消费组广播模式,按需使用。
- **组协调器:**运行在MQ服务上,管理一个消费组的元数据。消费者实例启动后回去找这个组协调器,加入进去由组协调器管理消费者,组协调器会给每个消费者分配分区列表,列表由kafka客户端维护,通过 getAssignedPartitions(); 获取,即可开始监听消费。
示例代码:
▼java复制代码// 消费者启动时,会从MQ服务端分配到一些分区,例如 partition-2 // getAssignedPartitions() 返回的是分配给【当前这一个消费者实例】的分区列表,而不是整个Topic的所有分区。 List<Partition> assignedPartitions = getAssignedPartitions(); // [partition-2] // 为每个分区启动一个独立的消费线程 for (Partition partition : assignedPartitions) { new Thread(() -> { while (true) { // 从 partition-2 中拉取一条消息 Message message = consumer.poll(partition); if (message != null) { // 同步处理消息,处理完再拉取下一条 processMessage(message); } } }).start(); }
RocketMQ 也是一样,由消息键 + 组内消费者实例独占队列实现。
消息积压
消息积压(Message Backlog)指的是生产者发送消息的速度超过了消费者处理消息的速度,导致大量消息在消息队列中堆积未被及时消费
消息积压会带来一系列负面影响:
- 业务延迟:消息处理延迟增加,导致业务操作(如下单、支付通知)响应变慢。
- 系统资源耗尽:堆积的消息占用大量 Broker 内存和磁盘资源,可能影响 Broker 性能甚至稳定性。
- 用户体验下降:用户感知到系统响应迟缓,操作卡顿,严重时可能导致服务不可用。
- 数据不一致风险:如果涉及事务或状态同步,消息积压可能导致数据在一段时间内不一致
消息积压通常由以下一个或多个因素导致:
- 消费能力不足:消费者数量不足、处理逻辑复杂(如涉及慢查询、外部调用)、单条消息处理时间长。
- 生产流量突增:如促销活动、秒杀场景导致瞬时流量激增,远超消费者处理能力。
- 消费者故障:消费者应用宕机、网络问题或 Full GC 等导致消费暂停。
- 资源配置不当:Topic 的队列数不足、消费者线程池配置不合理、QoS/Prefetch 参数设置不当等。
- Broker 资源瓶颈:磁盘 I/O、网络带宽或 CPU 资源受限,影响消息读写和分发
RabbitMQ 解决
RabbitMQ 解决消息积压的核心思路是提高消费能力和优化资源配置。思路是:监控 + 应急方案 + 预防
- 首先需要监控队列状态,及时发现积压。
-
- RabbitMQ 管理界面:查看队列的
**Ready**(待消费)、**Unacked**(已投递未确认)消息数 - 监控工具:集成 Prometheus + Grafana 监控队列深度、消费速率等指标,并设置告警(如积压超阈值)
- RabbitMQ 管理界面:查看队列的
- 应急处理方案
-
- **动态扩容消费者:**临时增加消费者实例数量。通过调整
**spring.rabbitmq.listener.simple.max-concurrency**来动态提高最大并发消费者数(需确保有足够资源) - 启用批量消费**:**调整
**prefetch**值:适当增加**prefetch**可以让消费者一次接收更多消息(但需权衡处理速度和内存占用) - 临时降级:如果资源有限,可暂时暂停非核心业务消费者或降低其线程数,优先保障核心业务消费
- **动态扩容消费者:**临时增加消费者实例数量。通过调整
- 预防消息积压
-
- **优化消费者逻辑:**异步化耗时操作;代码中检查慢查询、锁竞争等。
- **合理配置参数:**根据业务场景调整
**concurrency**(并发数)和**prefetch**。例如,IO密集型任务可设置较小的**prefetch**(如1-10),计算密集型可设置较大的**prefetch**(如100-200) - 交换机分离:为不同业务创建独立的队列和交换机,避免相互影响。
- 使用死信队列:对于处理失败的消息,配置死信队列(DLX)进行收集,避免无限重试阻塞队列;定期分析死信队列,修复问题根源。
RocketMQ 解决(推荐)
RocketMQ 凭借其高吞吐量和灵活的架构,在处理消息积压方面提供了强大的能力。
思路相同:监控 + 应急方案 + 预防(长期优化方案)。
- RocketMQ 提供了丰富的监控工具。
-
- RocketMQ 控制台:直接查看 Topic 的消息堆积量(
**diff**值)、消费者组的消费进度(Lag) - 命令行工具:使用
**mqadmin**命令(如**consumerProgress**、**topicStatus**)查询详细状态 - 监控集成:集成 Prometheus + Grafana,设置积压告警(如 Lag 超过 10 万条)
- RocketMQ 控制台:直接查看 Topic 的消息堆积量(
- 应急处理方案
-
- 动态扩容消费者:快速启动更多消费者实例。确保消费者组内的实例数与 Topic 的队列数匹配(例如,Topic 有 8 个队列,消费者组至少应有 8 个实例)以最大化并行度
- 优化消费线程:整消费者的线程池大小,增加消费并发度
- 启用批量消费:启用批量消费模式,减少网络开销。设置
**consumeMessageBatchMaxSize**(如每次拉取 32 条消息)
- 长期优化方案
-
- **增加队列数:**如果 Topic 队列数不足成为瓶颈,可通过
**updateTopic**命令增加队列数,提升消费并行度 - 优化 Broker 参数:将
**flushDiskType**设置为**ASYNC_FLUSH**(异步刷盘)以降低磁盘 I/O 压力 - **优化Broker参数:**调整
**sendMessageThreadPoolNums**和**pullMessageThreadPoolNums**等线程池参数,提升 Broker 处理能力 - 流量控制与降级:在紧急情况下,可对生产者进行限流,或对非核心消息启用降级策略(如丢弃部分日志类消息),防止积压进一步恶化
- **增加队列数:**如果 Topic 队列数不足成为瓶颈,可通过
小结
消息积压的解决思路
- 监控先行:建立完善的监控和告警机制是发现和解决消息积压的前提。
- 快速响应:制定应急预案,包括扩容、降级等流程,以便快速响应积压。
- 优化消费逻辑:大多数积压问题与消费逻辑有关,优化消费逻辑是根本。
- 合理规划:根据业务峰值进行容量规划,预留缓冲资源。
- 架构设计:对于高流量系统,考虑消息分片(Sharding)或多集群部署,提前分散压力。
事务消息
事务消息(Transaction Message)是一种特殊的消息类型,它能确保生产者发送消息与本地事务的原子性(要么都成功,要么都失败)。它通过 两阶段提交(2PC) 和事务状态回查机制,实现了分布式场景下的最终一致性。
大白话:生产者给消息队列发送一个半消息(不能被消费的事务消息),生产者在本地把其他事务操作都做了,全部成功,在给消息队列发送一个提交事务的消息,消息队列的半消息就变成普通消息了。
生命周期:
半消息 → 本地事务 → 提交/回滚 → 回查确认 → 消费投递。

事务消息主要为了解决分布式事务中跨服务操作的原子性问题。在微服务架构中,一个业务操作可能需要跨越多个服务(如订单服务、库存服务、支付服务),可能会发生数据不一致,事务消息的作用:
- 跨服务操作的一致性:确保核心业务操作(如订单创建)与多个下游操作(如库存扣减、积分变更、购物车清空)的结果完全一致
- 降低业务侵入:它将分布式事务的复杂性下沉到消息中间件,业务代码只需关注本地事务和简单的消息确认,开发更简单
举例:假设普通消息执行下面的操作失败了
▼java复制代码(1) 创建订单 → 成功 (2) 发送扣库存消息 → 失败
导致订单系统认为下单成功,但库存系统从没收到扣减请求 → 数据不一致
如果把这条普通消息设计成事务消息, 让上面的两步变成一个半事务操作:
- 要么订单创建成功 + 消息可靠发送;要么两者都不生效。
我们来看看RabbitMQ 和 RocketMQ 具体的落地实现
RabbitMQ 事务消息
RabbitMQ 的AMQP协议提供了原生的事务机制(**txSelect()**, **txCommit()**, **txRollback()**),通过阻塞式的方式保证强一致性
AMQP 协议: 是一个二进制的应用层协议,规定了消息队列系统的通信规范 ,规范中定义了 tx.select、tx.commit、tx.rollback 等命令 ,RabbitMQ 对这三个命令做了具体的实现。事务能力属于AMQP规范的一部分。
缺点:性能很差,很少实际使用RabbitMQ 的事务能力,需要事务能力技术选型应该选择 RocketMQ.
特点:
- 强一致性:通过协议层面保证,可靠性高。
- 阻塞式:事务提交是同步阻塞的,性能较低,吞吐量受限
- 适用场景:对一致性要求极高的场景,如金融交易。
工作流程:
- 开启事务:通过
**channel.txSelect()**将信道设置为事务模式。 - 发送消息:在事务中发送一条或多条消息。
- 提交或回滚:根据业务逻辑执行
**txCommit()**提交事务,或**txRollback()**回滚事务

RocketMQ 事务消息(异步最终一致)
RocketMQ 事务消息是高性能、易用的解决方案,采用异步化设计,通过半消息和状态回查机制实现最终一致性。
核心概念:
- 半消息:预先发送的消息,对消费者不可见,状态为“暂不能投递”。
- 本地事务:生产者执行的业务逻辑(如数据库操作)。
- 状态回查:Broker 定期回查生产者,确认本地事务的最终状态
工作流程:
- 发送半消息:生产者发送半消息到 Broker,Broker 持久化后返回 ACK 。
- 执行本地事务:生产者执行本地事务(如更新订单状态)。
- 提交或回滚:根据本地事务结果,向 Broker 发送
**COMMIT**或**ROLLBACK**。 - 状态回查:若 Broker 未收到确认,会定期回查生产者 。
特点:
- 最终一致性:不要求强一致性,性能高,适合高并发场景。
- 异步化:非阻塞,吞吐量高。
- 易用性:对业务代码的侵入性相对较低。

事务消息是分布式事务的一种实现方式。
分布式事务
指一个事务的参与者、涉及的资源服务器分布在不同的物理节点上。它需要保证跨多个独立资源(数据库、服务)的一系列操作,要么全部成功执行,要么全部失败回滚,从而维护数据的一致性。
分布式事务主要解决在分布式系统中,由于网络分区、节点故障、并发操作等带来的数据一致性问题。
实现分布式事务的方案有很多,主流的可以分为两大类:强一致性方案和最终一致性方案。
强一致性方案(同步阻塞)
这类方案追求所有节点在任何时刻的数据状态都是完全一致的。
-
两阶段提交(2PC - Two-Phase Commit)
-
- 角色:协调者、参与者。
- 第一阶段(准备阶段/投票阶段):协调者向所有参与者发送“prepare”请求,询问是否可以执行事务。参与者执行事务操作,但不提交,并锁定资源,然后向协调者返回“YES”或“NO”。
- 第二阶段(提交/回滚阶段):
-
-
- 如果所有参与者都返回“YES”,协调者发送“commit”请求,所有参与者正式提交事务。
- 如果任何一个参与者返回“NO”或超时,协调者发送“rollback”请求,所有参与者回滚事务。
-
-
- 缺点:同步阻塞、协调者单点故障、数据锁定时间长、性能差。
两阶段提交是分布式事务的模型,MySQL 的事务是经典实现,三个主流的消息队列没有对应实现。
2. 最终一致性方案(异步非阻塞)
这类方案允许系统在某一时刻数据不一致,但通过补偿机制,最终会达到一致状态。这是目前微服务架构下的主流选择。
TCC(Try-Confirm-Cancel)
TCC 是一种分布式事务处理模型或设计模式。它定义了一套标准的业务流程规范,将一个完整的业务逻辑按阶段拆分为三个方法:**Try**, **Confirm**, **Cancel**。任何遵循这套规范的实现,都可以称为 TCC 方案。
实现:
TCC 的核心思想是把一个大的分布式事务拆分成多个独立的本地事务,每个本地事务都有对应的三个操作。
核心角色
- 事务发起者:主业务服务,负责编排整个 TCC 事务流程。
- 事务参与者:被调用的其他微服务,需要实现 TCC 的三个接口。
- 事务协调器:通常由框架(如 Seata)承担,负责记录事务全局状态,并驱动 Confirm 或 Cancel 阶段的调用。
阶段详解
- Try 阶段(资源检查与预留)
-
- 目的:检查业务资源是否可用,并预留业务资源。
- 操作:这个阶段不执行真正的业务。例如,在支付场景中,
**Try**不是真的扣钱,而是冻结账户中的资金;在库存场景中,不是真的减库存,而是预占库存。 - 关键点:Try 操作需要具备幂等性,因为可能会被重试。
- Confirm 阶段(确认执行业务)
-
- 目的:当所有参与者(多个参与的微服务)
**Try**阶段都成功后,执行真正的业务操作。 - 操作:在支付场景中,
**Confirm**才是真正地扣减被冻结的资金;在库存场景中,才是真正地减去预占的库存。 - 关键点:Confirm 操作也必须具备幂等性。因为网络问题可能导致 Confirm 调用失败,协调器会重试,必须保证多次调用结果与一次调用相同。
- 目的:当所有参与者(多个参与的微服务)
- Cancel 阶段(取消预留资源)
-
- 目的:如果任何一个参与者(多个参与的微服务)的
**Try**阶段失败,则调用所有已成功**Try**的参与者的**Cancel**操作,回滚资源。 - 操作:在支付场景中,
**Cancel**是解冻之前冻结的资金;在库存场景中,是释放预占的库存。 - 关键点:Cancel 操作同样必须具备幂等性。
- 目的:如果任何一个参与者(多个参与的微服务)的
落地实现:
- Seata:阿里巴巴开源的分布式事务解决方案,提供了对 TCC 模式的完整支持。它扮演了一个事务协调者 的角色,帮助你管理 TCC 事务的生命周期,包括调用 Try、Confirm 或 Cancel 阶段。这是目前国内 Java 生态中最流行的选择。
- Hmily:一款高性能的分布式事务 TCC 服务框架,支持多种 RPC 框架,如 Dubbo、Spring Cloud 等。
- ByteTCC:一个基于 Java 的 TCC 型事务管理器,兼容 JTA(Java Transaction API)。
Saga模式
Saga 模式把一个 长事务(跨多个服务/多个数据库)拆分成一系列 本地事务(Local Transaction)。 每个本地事务在自己的服务/数据库中执行、提交。
如果某个本地事务失败了,那么之前已经成功的本地事务必须通过 补偿事务(Compensating Transaction) 撤销自己的效果。
通常适用于多个服务各自拥有数据库、且需要跨服务的一致性保障,但不能用传统分布式事务(如 2PC)或者不愿承担其性能成本。
落地技术/实现方式包括:
- 业务逻辑层面定义补偿接口
在每个本地事务步骤中,同时定义一个“正向操作(Forward Transaction)”和一个“补偿操作(Compensating Transaction)”。例如:库存服务有 reserveStock()(扣库存)与 releaseStock()(释放库存)作为补偿。
补偿操作通常作为异步任务或消息触发执行
- 事件驱动或消息队列触发
在 Saga 模式中,经常使用消息或事件来控制流程:当一个步骤成功后,发布事件触发下一个步骤;当失败时,发布“需要补偿”的事件,触发补偿事务。
- Orchestration 或 Choreography 模式控制
- 编排(Orchestration): 有一个中央协调器负责记录每一步状态,决定何时执行补偿。
- 编舞(Choreography): 各服务基于事件流触发自己的补偿逻辑,无中央协调器。
- 持久化 Saga 状态与日志
为了保证补偿事务可靠执行,系统通常会有一个 Saga 日志或状态存储,用来记录每一步是否成功、是否需要补偿。失败时,系统查日志执行补偿。
- 保证补偿操作幂等性
补偿事务可能会被重复触发,需要保证其幂等性(多次执行效果相同)
本地消息表
在一个服务内(本地数据库 + 该服务的消息发送逻辑)把 “业务操作 + 消息发送” 放在同一个本地事务中执行。 → 保证:如果业务表更新成功,消息也写入“消息表”;如果业务失败,消息也不会写入。
这样可避免“业务写成功但消息发送失败”导致的 “状态更新” + “通知消息”不一致问题。
不过很容易看出缺陷, 本地消息表的思路是确保每一个服务内的本地事务都是通过的,如果有一个不通过,就不会走到下游 ,但是上游的服务不会回滚, 它的特性是保证服务内原子,但不保证跨服务的原子回滚。
总结: “本地事务强”、但“跨服务事务弱”。
具体技术实现原理
下面是本地消息表模式的大致实现流程:
- 业务操作 + 消息插入在一个本地事务里 在服务 A 所使用的数据库内,进行如下操作:
▼plain复制代码BEGIN TRANSACTION 更新业务表(例如:创建订单) 插入消息表(例如:to_be_sent = true, message_body = {...}) COMMIT
这样保证了:业务数据变更和消息写入是原子的。
- 后台扫描或变更数据捕获(CDC)发送消息
-
- 一个后台任务或专门线程定期扫描
消息表中状态为to_be_sent的记录。 - 将消息发送到消息队列/中间件。
- 若发送成功,更新消息表为已发送状态(或者删除消息表记录)。
- 若失败,重试或记录失败。 这一机制保证“在业务成功且消息写入后再发送消息”。
- 一个后台任务或专门线程定期扫描
- 消费者消费 + 下一步服务执行 其他服务订阅队列消息,执行对应操作(例如库存扣减、支付通知等)。
- 幂等 + 重试机制
-
- 消息表发送过程可能重试多次;消费者也可能重复消费,因此系统需设计幂等逻辑。
- 消息可能延迟、重复,需要消费者做好幂等处理。
消费者‑队列 上的负载均衡
消费者 向 MQ 拉取数据还是MQ向消费者推送数据?
分情况讨论:
| 消息队列 | 消费模型 | 拉/推机制 | 说明 |
|---|---|---|---|
| Kafka | 消费者组(Pull) | ✅ 主动拉取(Pull) | 消费者定期向 broker 主动拉取消息 |
| RabbitMQ | 基于队列(Push) | ✅ 主动推送(Push) | Broker 主动将消息推送给消费者 |
| RocketMQ | 消费者组(Push + Pull混合) | ⚙️ 逻辑推送,本质拉取 | Broker 通知消费者有消息,消费者再发 Pull 请求拉取 |
1. RocketMQ 实现
RocketMQ 的消费者不是被动接收消息,而是 主动拉取(pull) 模式(即使是“推送式”消费者,也是在内部通过定时轮询拉取)。
负载均衡规则:RocketMQ 支持多种负载均衡规则,默认是平均分配。
例如:
- Topic:
OrderTopic - 有 8 个 MessageQueue
- Consumer Group:
OrderGroup - 启动了 4 个消费者实例
8个队列,4个消费者实例,负载均衡分配如下:
| Consumer | MessageQueues |
|---|---|
| C1 | MQ0, MQ1 |
| C2 | MQ2, MQ3 |
| C3 | MQ4, MQ5 |
| C4 | MQ6, MQ7 |
这就是平均分配,由于队列的消息可能会消费完毕,所以每隔20秒会执行一次重新平均分配。
或者其它触发逻辑:
- 消费者实例数量变化(新增或宕机)
- Topic 的 MessageQueue 数量变化(比如 Broker 扩容)
- 定时任务触发
- RabbitMQ 实现
在 RabbitMQ 中,由队列推送消息给消费者。
多消费者订阅同一个队列时,默认会轮询(round‑robin)分配消息给各个消费者。 并且是带补偿的轮询分发。
假设两个消费者 C1,C2订阅了队列 order_queue , 那么轮询规则是这样的:
| 消息 | 分发给谁 |
|---|---|
| msg1 | C1 |
| msg2 | C2 |
| msg3 | C1 |
| msg4 | C2 |
| … | … |
消费者消费完毕需要发回一个 ack 给队列,如果队列没有收到 ack , RabbitMQ 会把消息重新分配给别的消费者,确保这条消息被消费然后再消费下一条,这就是带补偿的轮询分发。
本文参考
https://rocketmq.apache.org/zh/docs/featureBehavior/03fifomessage/
https://rocketmq.apache.org/zh/docs/featureBehavior/04transactionmessage/
https://www.cnblogs.com/wunsiang/p/12765158.html
