MQ进阶 - 分布式下的挑战

上文讲了三大主流 MQ 的的功能和特性,那就像是工具箱,在实际场景中按需使用即可。

本文讲述的是使用 MQ 我们会遭遇的问题,相当于 BUG —— 使用 MQ 出现了不符合预期的结果,不解决不可上线。

引言 - 分布式的经典问题

在分布式场景下,给我们带来了高可用、易扩展和高容错的便利,也带来了分布式系统下不少的麻烦,增加了分布式开发的复杂性。

分布式的问题一图概览:

img

可知主要问题是三个根源引起:网络不可靠、节点独立性和无全局时钟。

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 的详细流程:

  1. 创建生产者: 应用程序创建一个 **DefaultMQProducer** 实例,并设置 NameServer 的地址。
  2. 启动生产者: 调用 **producer.start()** 方法。
    • 内部会启动多个服务线程。
    • 与 NameServer 建立长连接,定期(默认每30秒)从 NameServer 拉取最新的 Topic 路由信息(即该 Topic 对应的 Broker 地址、队列等)。
  1. 构建消息: 创建一个 **Message** 对象,包含 Topic、Body(消息体)和可选的 Tag/Keys。
  2. 发送消息: 调用 **producer.send(message)** 方法。
    • 选择队列: 生产者客户端会根据预设 (代码逻辑预设) 的负载均衡策略(如轮询),从获取到的路由信息中选择一个消息队列(Message Queue,位于某个 Broker 上)。
    • 网络发送: 将消息发送给目标 Broker。
  1. 等待 Broker ACK (关键步骤):
    • 默认情况下,**send()** 方法是同步阻塞的。它会一直等待,直到收到来自 Broker 的响应,或者直到超时。
    • Broker 接收到消息后,会将其写入 CommitLog(内存和磁盘),然后向生产者返回一个 **SendResult** 对象。
    • 这个 **SendResult** 就是生产端的 ACK。它包含了发送状态(**SendStatus**),如 **SEND_OK**(成功)、**FLUSH_DISK_TIMEOUT**(刷盘超时)、**FLUSH_SLAVE_TIMEOUT**(同步到从机超时)、**SLAVE_NOT_AVAILABLE**(从机不可用)等。
  1. 处理结果: 生产者根据 **SendResult** 的状态,可以判断消息是否成功发送,并进行相应的业务处理(如重试、记录日志等)。

img

RocketMQ 生产消息到 Brocker 的详细流程:

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

img

RabbitMQ

核心角色

  • Producer (生产者): 消息的发送方。负责将消息发送到交换机。
  • Consumer (消费者): 消息的接收方。负责从队列中获取消息并进行处理。
  • Broker (代理服务器): RabbitMQ 服务本身。负责接收、存储和路由消息。
  • Exchange (交换机): 接收来自生产者的消息,并根据路由键将消息推送到一个或多个队列。生产者从不直接将消息发送到队列。
  • Queue (队列): 存储消息的缓冲区,直到消费者将其取走。
  • Channel (信道): 位于 TCP 连接内部的虚拟连接。它是进行 AMQP 操作(如发布、消费、获取消息)的轻量级载体,避免了为每个操作都建立 TCP 连接的开销。

RabbitMQ 生产消息到队列的详细流程:

  1. 创建连接和信道: 应用程序与 Broker 建立一个 TCP 连接,并在其上创建一个信道。
  2. 启用发布者确认 (关键步骤): 在信道上调用 **channel.confirmSelect()**,将该信道设置为确认模式。这是开启生产端 ACK 的前提
  3. 声明交换机和队列: 为了保证消息能被正确路由,生产者通常会声明一个交换机和一个队列,并将它们绑定起来(这是一个幂等操作,重复声明不会出错)。
  4. 发送消息: 调用 **channel.basicPublish()** 方法,将消息发送到指定的交换机,并附带路由键。
  5. 等待 Broker ACK (关键步骤):
    • 消息发送后,生产者不会立即收到响应。Broker 接收到消息后,会将其路由到匹配的队列。
    • 一旦消息至少被一个队列接收并持久化(如果队列是持久化的),Broker 就会通过该信道向生产者发送一个确认帧,包含消息的 **delivery-tag**
    • 如果由于某些原因(例如找不到匹配的队列)消息无法被路由,Broker 会发送一个未确认帧
  1. 处理确认: 生产者可以监听信道的确认和未确认事件,根据收到的 **delivery-tag** 来判断哪条消息成功或失败,并进行相应的补偿业务处理(如记录日志、重发失败的消息等)。

img

RabbitMQ 消费流程:

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

img

消息丢失

消息丢失指的是:生产者发送了消息,但这条消息未能被成功存储或被消费者处理,最终在系统中“消失”,不再被传递或消费。消息丢失问题会严重影响数据一致性和业务可靠性

问题产生

消息丢失通常有以下几个环节可能发生:

  • 生产者发送失败:生产者向 Broker 发送消息时,可能因为网络中断、Broker 不可用、客户端超时而失败。如果生产者误以为发送成功但实际上没有,则消息丢失。
  • Broker 存储失败或确认失败:消息到达 Broker,但在持久化(写磁盘)或写入日志、刷盘、同步副本等过程中失败,或者 Broker 挂掉,还未完成持久化。
  • 消费者处理失败并且未重试/未记录:消费者接收消息后处理失败,且系统没有做重试或者死信机制、也没有记录该消息状态,导致消息被“吞掉”。
  • 配置或使用不当:例如队列或主题未配置持久化、生产者未启用确认机制、Broker 副本机制配置不当、网络故障未被捕捉等。

总结一句话:消息丢失的根因通常是网络、存储、确认机制不可靠,原因很多,无法全部概括。

问题解决

确保消息不丢失需要从生产者、Broker、消费者三个层面进行保障

  1. 生产者端
    • 使用发送确认机制:确保消息成功发送到Broker。
    • 配置重试机制:发送失败时自动重试
  1. Broker端
    • 消息持久化:将消息和队列设置为持久化,写入磁盘
    • 配置副本机制:通过主从复制或多副本保证数据冗余
  1. 消费者端
    • 手动提交offset:在消息处理完成后手动确认(程序员代码中用业务逻辑控制,不依赖消息中间件提供的自动提交),避免自动提交导致的问题
    • 消费幂等性设计:确保重复消费不会导致业务逻辑错误

但是这不能百分百保证消息丢失,所以下一个小节会讲消息补充,我们能接受极端情况下的丢失。

RabbitMQ 解决

RabbitMQ通过多层次确认机制持久化策略保障,其核心思想是在消息传递的每个关键节点设置确认点。

  1. 生产者确认机制
    • Publisher Confirms: 当生产者发送消息到RabbitMQ后,Broker会返回一个确认结果(ACK / NACK)
    • Publisher Returns : RabbitMQ 提供了 ReturnCallback 回调机制 ,开启mandatory=true 发送消息 时, 如果交换机找不到任何队列匹配,Broker 会 立即回调生产者,告诉你消息无法路由。 这样生产者就知道有消息没被队列接收,可以选择 重试、记录日志或发送到死信队列
  1. 持久化策略

为了保证消息不丢失,RabbitMQ不仅仅持久化了消息,还持久化了 Exchange 和 Queue , 做到了三层立体式持久化。

层级持久化对象含义作用
① Exchange 层交换机持久化Broker 重启后,交换机配置(名称、类型、绑定关系)仍然存在确保消息路由基础结构不丢
② Queue 层队列持久化Broker 重启后,队列仍然存在(包括绑定信息)确保消息能有地方落地
③ Message 层消息持久化消息内容写入磁盘(而不是只在内存)确保消息数据不丢失
  1. 消费者确认机制

消息确认是消费者向Broker发送确认,表示这条消息处理完了, Broker 收到 ack 后,才会把这条消息从队列里删除。

**RabbitMQ 支持自动确认,**缺点是容易丢失消息,因为自动确认机制是在消费者真实消费数据之前发送 ACK 。自动确认吞吐量高。

手动确认:可以自己写代码,灵活控制发送ACK的时机,更稳定。

RocketMQ 解决

  1. 生产者可靠性机制
  • 同步发送:使用同步发送方式,等待Broker返回确认结果,例如**SendResult sendResult = producer.send(msg)**。通过检查**sendResult.getSendStatus()**判断是否发送成功,失败时可进行重试
  • 事务消息:对一致性要求极高的场景,使用事务消息。它通过两阶段提交(2PC)确保本地事务与消息发送的原子性
  1. Broker端可靠性机制
  • 同步刷盘(SYNC_FLUSH):消息写入内存后立即刷盘,确保数据落盘,防止因Broker突然宕机导致内存中未刷盘的消息丢失
  • 异步刷盘(ASYNC_FLUSH):消息写入内存后异步刷盘,性能更好
  • 主从复制:配置主从Broker。将**brokerRole**设置为**SYNC_MASTER**(同步复制)。Master接收到消息后,会同步等待Slave复制成功,才向生产者返回ACK。即使Master宕机,数据也在Slave上有备份
  1. 消费者可靠性机制
  • 消费成功后才提交Offset:消费者在业务逻辑执行成功后,再向Broker提交消费进度。如果处理过程中消费者宕机,消息会重新投递,不会丢失
  • 避免异步消费:RocketMQ官方建议不要在消费端使用异步处理消息,以免因未处理完就提前提交Offset导致消息丢失

img

消息丢失问题不能根除(ack丢失,刷盘崩溃等), 接受极端情况丢失,但用补偿机制兜底

消息补偿

消息补偿是一种事后补救机制,用于在消息传递或处理失败时,通过重新发送或执行补偿操作,确保业务数据最终达到一致状态,它通常用于处理分布式事务中的异常情况,特别是在采用最终一致性模型的系统中。

具体来说,生产者发送消息失败、消费者消费失败、消息超时未确认、事务状态未知会触发消息补偿。

rabbitMQ 消息补偿机制

  1. 生产者确认机制

生产者发送消息时,可开启 Confirm 模式。Broker 会异步确认消息是否到达交换机,会出现两种情况:

  • ACK:消息成功到达交换机。
  • NACK:消息未到达交换机。生产者可在回调中捕获 NACK,执行重新发送记录日志等补偿操作
  1. 死信队列

当消息成为“死信”(如消费者拒绝且不重新入队、消息过期、队列满)时,若配置了死信交换机(DLX),消息会被路由到死信队列。进入死信队列的消息,可在程序中写逻辑对死信队列进行重试、重发(让生成者再发一次)、记录日志、人工介入。

  1. 消息重试与消费者确认

消费者处理失败时,可拒绝消息并选择不重新入队**requeue=false**),让消息进入死信队列,避免无限重试导致阻塞。同时,应确保业务逻辑的幂等性,以应对可能的消息重复

RocketMQ 消息补偿机制

RocketMQ 提供了更强大的事务消息机制,并依赖定时任务状态回查来实现补偿。

  1. 事务消息

补偿机制(事务状态回查):若Broker长时间未收到确认,会反向回查生产者的事务状态,生产者需实现回查接口,根据本地事务状态决定提交或回滚。事务回查触发消息补偿。

  1. 定时任务与消息重试

生产者:发送失败时可自动重试,自带消息补偿。

消费者:消费失败时,消息会自动重试。超过最大重试次数后,消息会进入**死信队列,**可监控死信队列并进行补偿。

  1. 本地消息表

对于未使用事务消息的场景,可采用本地消息表方案

本地消息表: 在生成者端建立一个消息表,用来把消息存储在业务数据库中,和业务操作放在同一个本地事务里提交。 这样就能确保业务操作成功的同时,消息也一定成功发送,避免中间的缝隙发生其它意外导致业务操作成功,消息发送失败的情况。

本地消息表是一种补偿机制, 先落地保存“要发的消息”,等网络恢复、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)指的是消息按照生产者发送的先后顺序被消费者接收和处理。在消息队列中,这通常分为两种类型

  1. 全局有序:整个主题(Topic)中的所有消息,都严格按照生产者发送的先后顺序进行消费。例如,发送顺序是 M1, M2, M3,消费时也必须是 M1, M2, M3。
  2. 分区有序(消息键分区):只需要保证具有某个相同特征(如同一个订单ID、同一个用户ID)的消息能够被顺序消费即可。例如,订单 **A** 的消息 **A1, A2, A3** 顺序消费,订单 **B** 的消息 **B1, B2, B3** 也顺序消费,但 **A1****B1** 谁先消费无所谓。

消息顺序的重要性

操作的顺序至关重要,乱序会导致业务逻辑错误或数据不一致。

比如订单的状态流转:创建订单 → 支付订单 → 发货 → 确认收货,一个状态是一条消息,上游执行成功会给下游发送一个消息,如果消息不会顺序消费,那么会发生数据不一致,业务混乱。

在这个场景中,每一个状态肯定是顺序发给队列的,只要确保同一个订单的消息进入同一个队列,那么在队列中的消息就是物理有序,RocketMQ 默认是消费组内单实例的单线程消费,不会造成乱序的问题。

作为开发者,我们只要保证同一个订单的消息发送给固定的一个队列就好,这就要用到消息键。

消息乱序的原因

消息队列(如 RabbitMQ、Kafka、RocketMQ、Redis Stream)都是并发消费模型。 如果你启用了多个消费者或多线程消费,一个消息队列中不同消息可能被不同线程同时处理

实现消息有序

我们知道有消息顺序:全局有序和分区有序,其中分区有序是重点难点,使用的也是比较多的,因为分区有序性能高,但是复杂一点。

分区有序使用消息键模式实现,接下来看各个MQ怎么处理。

kafka 与 Rocket MQ 实现方案相同

消息键:在生成消息的时候指定 **一个能区分业务功能、能够分组业务的字段 ,**字段可以是任何类型,会经过哈希计算,对Brocker的总队列取模,相同的业务便投入到固定队列。

以kafka为语境进行讲解。

原理详解:

  1. 生产者将消息路由到固定分区(也就是队列)

    1. 定义业务键:为每条消息设定一个能够标识其业务分组的键。例如,处理订单消息,**order_id** 就是绝佳的业务键。
    2. 根据键计算分区:生产者在发送消息时,不是随机选择分区,而是根据业务键进行哈希计算,然后对分区总数取模,从而确定该消息应该发往哪个分区。
    3. 示例代码:
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)。

  1. 消费者端单线程消费分区

    1. **天生支持单线程消费:**消息队列的消费者组模型天然支持这一点。一个分区在同一时间只能被消费者组里的一个消费者实例消费。
    2. 多个消费组可以订阅一个topic, 那么多个消息组都有自己的 offset , 一条消息会被多个消费组消费,这是跨消费组广播模式,按需使用。
    3. **组协调器:**运行在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 解决消息积压的核心思路是提高消费能力优化资源配置。思路是:监控 + 应急方案 + 预防

  1. 首先需要监控队列状态,及时发现积压。
    • RabbitMQ 管理界面:查看队列的 **Ready**(待消费)、**Unacked**(已投递未确认)消息数
    • 监控工具:集成 Prometheus + Grafana 监控队列深度、消费速率等指标,并设置告警(如积压超阈值)
  1. 应急处理方案
    • **动态扩容消费者:**临时增加消费者实例数量。通过调整 **spring.rabbitmq.listener.simple.max-concurrency** 来动态提高最大并发消费者数(需确保有足够资源)
    • 启用批量消费**:**调整 **prefetch** 值:适当增加 **prefetch** 可以让消费者一次接收更多消息(但需权衡处理速度和内存占用)
    • 临时降级:如果资源有限,可暂时暂停非核心业务消费者或降低其线程数,优先保障核心业务消费
  1. 预防消息积压
    • **优化消费者逻辑:**异步化耗时操作;代码中检查慢查询、锁竞争等。
    • **合理配置参数:**根据业务场景调整 **concurrency**(并发数)和 **prefetch**。例如,IO密集型任务可设置较小的 **prefetch**(如1-10),计算密集型可设置较大的 **prefetch**(如100-200)
    • 交换机分离:为不同业务创建独立的队列和交换机,避免相互影响。
    • 使用死信队列:对于处理失败的消息,配置死信队列(DLX)进行收集,避免无限重试阻塞队列;定期分析死信队列,修复问题根源。

RocketMQ 解决(推荐)

RocketMQ 凭借其高吞吐量灵活的架构,在处理消息积压方面提供了强大的能力。

思路相同:监控 + 应急方案 + 预防(长期优化方案)。

  1. RocketMQ 提供了丰富的监控工具。
    • RocketMQ 控制台:直接查看 Topic 的消息堆积量(**diff** 值)、消费者组的消费进度(Lag)
    • 命令行工具:使用 **mqadmin** 命令(如 **consumerProgress****topicStatus**)查询详细状态
    • 监控集成:集成 Prometheus + Grafana,设置积压告警(如 Lag 超过 10 万条)
  1. 应急处理方案
    • 动态扩容消费者:快速启动更多消费者实例。确保消费者组内的实例数与 Topic 的队列数匹配(例如,Topic 有 8 个队列,消费者组至少应有 8 个实例)以最大化并行度
    • 优化消费线程:整消费者的线程池大小,增加消费并发度
    • 启用批量消费:启用批量消费模式,减少网络开销。设置 **consumeMessageBatchMaxSize**(如每次拉取 32 条消息)
  1. 长期优化方案
    • **增加队列数:**如果 Topic 队列数不足成为瓶颈,可通过 **updateTopic** 命令增加队列数,提升消费并行度
    • 优化 Broker 参数:将 **flushDiskType** 设置为 **ASYNC_FLUSH**(异步刷盘)以降低磁盘 I/O 压力
    • **优化Broker参数:**调整 **sendMessageThreadPoolNums****pullMessageThreadPoolNums** 等线程池参数,提升 Broker 处理能力
    • 流量控制与降级:在紧急情况下,可对生产者进行限流,或对非核心消息启用降级策略(如丢弃部分日志类消息),防止积压进一步恶化

小结

消息积压的解决思路

  1. 监控先行:建立完善的监控和告警机制是发现和解决消息积压的前提。
  2. 快速响应:制定应急预案,包括扩容、降级等流程,以便快速响应积压。
  3. 优化消费逻辑:大多数积压问题与消费逻辑有关,优化消费逻辑是根本。
  4. 合理规划:根据业务峰值进行容量规划,预留缓冲资源。
  5. 架构设计:对于高流量系统,考虑消息分片(Sharding)或多集群部署,提前分散压力。

事务消息

事务消息(Transaction Message)是一种特殊的消息类型,它能确保生产者发送消息与本地事务的原子性(要么都成功,要么都失败)。它通过 两阶段提交(2PC)事务状态回查机制,实现了分布式场景下的最终一致性。

大白话:生产者给消息队列发送一个半消息(不能被消费的事务消息),生产者在本地把其他事务操作都做了,全部成功,在给消息队列发送一个提交事务的消息,消息队列的半消息就变成普通消息了。

生命周期:

半消息 → 本地事务 → 提交/回滚 → 回查确认 → 消费投递。

img

事务消息主要为了解决分布式事务跨服务操作的原子性问题。在微服务架构中,一个业务操作可能需要跨越多个服务(如订单服务、库存服务、支付服务),可能会发生数据不一致,事务消息的作用:

  • 跨服务操作的一致性:确保核心业务操作(如订单创建)与多个下游操作(如库存扣减、积分变更、购物车清空)的结果完全一致
  • 降低业务侵入:它将分布式事务的复杂性下沉到消息中间件,业务代码只需关注本地事务和简单的消息确认,开发更简单

举例:假设普通消息执行下面的操作失败了

java
复制代码
(1) 创建订单 → 成功 (2) 发送扣库存消息 → 失败

导致订单系统认为下单成功,但库存系统从没收到扣减请求 → 数据不一致

如果把这条普通消息设计成事务消息, 让上面的两步变成一个半事务操作:

  • 要么订单创建成功 + 消息可靠发送;要么两者都不生效。

我们来看看RabbitMQ 和 RocketMQ 具体的落地实现

RabbitMQ 事务消息

RabbitMQ 的AMQP协议提供了原生的事务机制(**txSelect()**, **txCommit()**, **txRollback()**),通过阻塞式的方式保证强一致性

AMQP 协议: 是一个二进制的应用层协议,规定了消息队列系统的通信规范 ,规范中定义了 tx.selecttx.committx.rollback 等命令 ,RabbitMQ 对这三个命令做了具体的实现。事务能力属于AMQP规范的一部分。

缺点:性能很差,很少实际使用RabbitMQ 的事务能力,需要事务能力技术选型应该选择 RocketMQ.

特点

  • 强一致性:通过协议层面保证,可靠性高。
  • 阻塞式:事务提交是同步阻塞的,性能较低,吞吐量受限
  • 适用场景:对一致性要求极高的场景,如金融交易。

工作流程

  1. 开启事务:通过 **channel.txSelect()** 将信道设置为事务模式。
  2. 发送消息:在事务中发送一条或多条消息。
  3. 提交或回滚:根据业务逻辑执行 **txCommit()** 提交事务,或 **txRollback()** 回滚事务

img

RocketMQ 事务消息(异步最终一致)

RocketMQ 事务消息是高性能、易用的解决方案,采用异步化设计,通过半消息状态回查机制实现最终一致性。

核心概念

  • 半消息:预先发送的消息,对消费者不可见,状态为“暂不能投递”。
  • 本地事务:生产者执行的业务逻辑(如数据库操作)。
  • 状态回查:Broker 定期回查生产者,确认本地事务的最终状态

工作流程

  1. 发送半消息:生产者发送半消息到 Broker,Broker 持久化后返回 ACK 。
  2. 执行本地事务:生产者执行本地事务(如更新订单状态)。
  3. 提交或回滚:根据本地事务结果,向 Broker 发送 **COMMIT****ROLLBACK**
  4. 状态回查:若 Broker 未收到确认,会定期回查生产者 。

特点

  • 最终一致性:不要求强一致性,性能高,适合高并发场景。
  • 异步化:非阻塞,吞吐量高。
  • 易用性:对业务代码的侵入性相对较低。

img

事务消息是分布式事务的一种实现方式。

分布式事务

指一个事务的参与者、涉及的资源服务器分布在不同的物理节点上。它需要保证跨多个独立资源(数据库、服务)的一系列操作,要么全部成功执行,要么全部失败回滚,从而维护数据的一致性。

分布式事务主要解决在分布式系统中,由于网络分区、节点故障、并发操作等带来的数据一致性问题。

实现分布式事务的方案有很多,主流的可以分为两大类:强一致性方案最终一致性方案

强一致性方案(同步阻塞)

这类方案追求所有节点在任何时刻的数据状态都是完全一致的。

  • 两阶段提交(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 阶段的调用。

阶段详解

  1. Try 阶段(资源检查与预留)
    • 目的:检查业务资源是否可用,并预留业务资源。
    • 操作:这个阶段不执行真正的业务。例如,在支付场景中,**Try** 不是真的扣钱,而是冻结账户中的资金;在库存场景中,不是真的减库存,而是预占库存。
    • 关键点:Try 操作需要具备幂等性,因为可能会被重试。
  1. Confirm 阶段(确认执行业务)
    • 目的:当所有参与者(多个参与的微服务) **Try** 阶段都成功后,执行真正的业务操作。
    • 操作:在支付场景中,**Confirm** 才是真正地扣减被冻结的资金;在库存场景中,才是真正地减去预占的库存。
    • 关键点:Confirm 操作也必须具备幂等性。因为网络问题可能导致 Confirm 调用失败,协调器会重试,必须保证多次调用结果与一次调用相同。
  1. 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)或者不愿承担其性能成本。

落地技术/实现方式包括:

  1. 业务逻辑层面定义补偿接口

在每个本地事务步骤中,同时定义一个“正向操作(Forward Transaction)”和一个“补偿操作(Compensating Transaction)”。例如:库存服务有 reserveStock()(扣库存)与 releaseStock()(释放库存)作为补偿。

补偿操作通常作为异步任务或消息触发执行

  1. 事件驱动或消息队列触发

在 Saga 模式中,经常使用消息或事件来控制流程:当一个步骤成功后,发布事件触发下一个步骤;当失败时,发布“需要补偿”的事件,触发补偿事务。

  1. Orchestration 或 Choreography 模式控制
  • 编排(Orchestration): 有一个中央协调器负责记录每一步状态,决定何时执行补偿。
  • 编舞(Choreography): 各服务基于事件流触发自己的补偿逻辑,无中央协调器。
  1. 持久化 Saga 状态与日志

为了保证补偿事务可靠执行,系统通常会有一个 Saga 日志或状态存储,用来记录每一步是否成功、是否需要补偿。失败时,系统查日志执行补偿。

  1. 保证补偿操作幂等性

补偿事务可能会被重复触发,需要保证其幂等性(多次执行效果相同)


本地消息表

在一个服务内(本地数据库 + 该服务的消息发送逻辑)把 “业务操作 + 消息发送” 放在同一个本地事务中执行。 → 保证:如果业务表更新成功,消息也写入“消息表”;如果业务失败,消息也不会写入。

这样可避免“业务写成功但消息发送失败”导致的 “状态更新” + “通知消息”不一致问题。

不过很容易看出缺陷, 本地消息表的思路是确保每一个服务内的本地事务都是通过的,如果有一个不通过,就不会走到下游 ,但是上游的服务不会回滚, 它的特性是保证服务内原子,但不保证跨服务的原子回滚

总结: “本地事务强”、但“跨服务事务弱”。

具体技术实现原理

下面是本地消息表模式的大致实现流程:

  1. 业务操作 + 消息插入在一个本地事务里 在服务 A 所使用的数据库内,进行如下操作:
plain
复制代码
BEGIN TRANSACTION 更新业务表(例如:创建订单) 插入消息表(例如:to_be_sent = true, message_body = {...}) COMMIT

这样保证了:业务数据变更和消息写入是原子的。

  1. 后台扫描或变更数据捕获(CDC)发送消息
    • 一个后台任务或专门线程定期扫描 消息表 中状态为 to_be_sent 的记录。
    • 将消息发送到消息队列/中间件。
    • 若发送成功,更新消息表为已发送状态(或者删除消息表记录)。
    • 若失败,重试或记录失败。 这一机制保证“在业务成功且消息写入后再发送消息”。
  1. 消费者消费 + 下一步服务执行 其他服务订阅队列消息,执行对应操作(例如库存扣减、支付通知等)。
  2. 幂等 + 重试机制
    • 消息表发送过程可能重试多次;消费者也可能重复消费,因此系统需设计幂等逻辑。
    • 消息可能延迟、重复,需要消费者做好幂等处理。

消费者‑队列 上的负载均衡

消费者 向 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个消费者实例,负载均衡分配如下:

ConsumerMessageQueues
C1MQ0, MQ1
C2MQ2, MQ3
C3MQ4, MQ5
C4MQ6, MQ7

这就是平均分配,由于队列的消息可能会消费完毕,所以每隔20秒会执行一次重新平均分配。

或者其它触发逻辑:

  • 消费者实例数量变化(新增或宕机)
  • Topic 的 MessageQueue 数量变化(比如 Broker 扩容)
  • 定时任务触发
  1. RabbitMQ 实现

在 RabbitMQ 中,由队列推送消息给消费者。

多消费者订阅同一个队列时,默认会轮询(round‑robin)分配消息给各个消费者。 并且是带补偿的轮询分发

假设两个消费者 C1,C2订阅了队列 order_queue , 那么轮询规则是这样的:

消息分发给谁
msg1C1
msg2C2
msg3C1
msg4C2

消费者消费完毕需要发回一个 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

https://juejin.cn/post/7491971657202073654

https://www.cnblogs.com/RunningSnails/p/17373107.html

0个评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
下载 APP