消息队列

MQ 的核心价值、可靠性问题与常见机制梳理

Posted by Ekko on August 19, 2020

内容围绕消息队列在分布式系统中的角色展开,重点梳理解耦、异步、削峰、事件驱动,以及重复消费、顺序消费、消息丢失和分布式事务等核心问题。

文章偏向基础总览,关注常见概念之间的联系,以及工程实践中最容易踩到的边界条件。

参考资料:

官方文档:Jakarta Messaging 3.1RabbitMQ AMQP 0-9-1 ModelRabbitMQ ExchangesApache RocketMQ Ordered Message

实践资料:掘金消息队列3y三太子敖丙JavaGuide芋道源码

[TOC]


什么是消息队列

MQ 全称 Message Queue,消息队列(MQ)是一种应用程序对应用程序的通信方式。应用程序通过写入和读取消息来完成通信,而不需要像同步 RPC 那样要求调用方和被调用方在同一时刻直接建立强依赖关系。消息传递强调的是通过消息载体交换数据,排队强调的是把待处理的数据暂存在中间层,由后续消费者按约定的方式继续处理。

把消息队列比作是一个存放消息的容器,当需要使用消息的时候可以取出消息供自己使用。消息队列是分布式系统中重要的组件,使用消息队列主要是为了通过 异步 处理提高系统性能和削峰降低系统耦合性。目前使用较多的消息队列有 ActiveMQ、RabbitMQ、Kafka、RocketMQ

需要注意,消息队列并不只是一个简单的 FIFO 队列抽象。Kafka 更接近可追加的分区日志,RabbitMQ 通过 exchange 和 queue 完成路由,RocketMQ 则使用 topic 和 queue 组织消息。因此在工程实践里,更重要的问题往往不是“有没有一个队列”,而是消息如何被路由、持久化、确认、重试,以及在什么范围内保证顺序。

除了上面说的消息消费顺序的问题,使用消息队列,我们还要考虑如何保证消息不被重复消费?如何保证消息的可靠性传输(如何处理消息丢失的问题)?……等等问题。所以说使用消息队列也不是十全十美的,使用它也会让系统可用性降低、复杂度提高,另外需要我们保障一致性等问题

优点: 异步、削峰、解耦 缺点: 重复消费、消息丢失、顺序消费


解耦、异步、削峰

解耦:

传统模式:

消息队列之解耦传统模式.png

缺点: 系统间耦合性太强,如上图所示,系统A在代码中直接调用系统B和系统C的代码,如果将来D系统接入,系统A还需要修改代码,过于麻烦

中间件模式:

消息队列之解耦中间件模式.png

优点: 将消息写入消息队列,需要消息的系统自己从消息队列中订阅,从而系统A不需要做任何修改

如果模块之间不存在直接调用,那么新增模块或者修改模块就对其他模块影响较小,这样系统的可扩展性无疑更好一些

最常见的事件驱动架构类似生产者消费者模式,在大型网站中通常用利用消息队列实现事件驱动结构。如下图所示

消息队列实现事件驱动.png

消息队列是利用发布-订阅模式工作,消息发送者(生产者)发布消息,一个或多个消息接受者(消费者)订阅消息。 从上图可以看到消息发送者(生产者)和消息接受者(消费者)之间没有直接耦合,消息发送者将消息发送至分布式消息队列即结束对消息的处理,消息接受者从分布式消息队列获取该消息后进行后续处理,并不需要知道该消息从何而来。对新增业务,只要对该类消息感兴趣,即可订阅该消息,对原有系统和业务没有任何影响,从而实现网站业务的可扩展性设计

消息接受者对消息进行过滤、处理、包装后,构造成一个新的消息类型,将消息继续发送出去,等待其他消息接受者订阅该消息。因此基于事件(消息对象)驱动的业务架构可以是一系列流程

为了降低消息丢失风险,工程上通常不会把可靠性寄托在单一组件上,而是会同时治理三段链路:生产者侧通过发送确认、事务消息或本地消息表保证“该发出去的消息确实发出去了”;Broker 侧通过持久化、刷盘和副本机制提升存储可靠性;消费者侧通过确认机制、重试、死信队列和幂等控制保证消息最终被正确处理。

消息队列并不只对应发布-订阅一种工作方式。JMS 规范从 Java API 视角抽象出 Queue 和 Topic 两种消息域;而 AMQP 这类协议则进一步定义了 broker、exchange、queue、binding、ack 等更底层的消息传递模型。二者描述的不是同一个层次的问题。


异步:

传统模式:

消息队列之异步传统模式.png

缺点: 一些非必要的业务逻辑以同步的方式运行,太耗费时间

消息队列之异步中间件模式.png

优点: 将消息写入消息队列,非必要的业务逻辑以异步的方式运行,加快响应速度

在不使用消息队列服务器的时候,用户的请求数据直接写入数据库,在高并发的情况下数据库压力剧增,使得响应速度变慢。但是在使用消息队列之后,用户的请求数据发送给消息队列之后立即返回,再由消息队列的消费者进程从消息队列中获取数据,异步写入数据库。由于消息队列服务器处理速度快于数据库(消息队列也比数据库有更好的伸缩性),因此响应速度得到大幅改善

通过以上分析我们可以得出消息队列具有很好的削峰作用的功能——即通过异步处理,将短时间高并发产生的事务消息存储在消息队列中,从而削平高峰期的并发事务

举例: 在电子商务一些秒杀、促销活动中,合理使用消息队列可以有效抵御促销活动刚开始大量订单涌入对系统的冲击。平时流量很低,但是你要做秒杀活动 00 :00 的时候流量疯狂怼进来,服务器,Redis,MySQL 各自的承受能力都不一样,全部流量照单全收肯定有问题,直接就打挂了 DB

因为用户请求数据写入消息队列之后就立即返回给用户了,但是请求数据在后续的业务校验、写数据库等操作中可能失败。因此使用消息队列进行异步处理之后,需要适当修改业务流程进行配合,比如用户在提交订单之后,订单数据写入消息队列,不能立即返回用户订单提交成功,需要在消息队列的订单消费者进程真正处理完该订单之后,甚至出库后,再通过电子邮件或短信通知用户订单成功,以免交易纠纷。这就类似我们平时手机订火车票和电影票


削峰:

消息队列之削峰传统模式.png

缺点: 并发量大的时候,所有的请求直接怼到数据库,造成数据库连接异常

消息队列之削峰中间件模式.png

优点: 系统 A 慢慢的按照数据库能处理的并发量,从消息队列中慢慢拉取消息。在生产中,这个短暂的高峰期积压是允许的

事件驱动架构

事件驱动架构(Event Driven Architecture,EDA)一个事件驱动框架(EDA)定义了一个设计和实现一个应用系统的方法学,在这个系统里事件可传输于松散耦合的组件和服务之间

一个事件驱动系统典型地由事件消费者和事件产生者组成。事件消费者向事件管理器订阅事件,事件产生者向事件管理器发布事件。当事件管理器从事件产生者那接收到一个事件时,事件管理把这个事件转送给相应的事件消费者。如果这个事件消费者是不可用的,事件管理者将保留这个事件,一段间隔之后再次转送该事件消费者。这种事件传送方法在基于消息的系统里就是:储存(store)和转送(forward)


使用消息队列带来的一些问题

系统可用性降低: 系统可用性在某种程度上降低,在加入 MQ 之前,不用考虑消息丢失或者说 MQ 挂掉等等的情况,但是,引入 MQ 之后就需要去考虑

系统复杂性提高: 加入 MQ 之后,你需要保证消息没有被重复消费、处理消息丢失的情况、保证消息传递的顺序性等等问题

数据一致性问题:

上面讲了消息队列可以实现异步,消息队列带来的异步确实可以提高系统响应速度。但是,万一消息的真正消费者并没有正确消费消息怎么办?这样就会导致数据不一致的情况了

这个其实是分布式服务本身就存在的一个问题,不仅仅是消息队列的问题,但是放在这里说是因为用了消息队列这个问题会暴露得比较严重一点

下单的服务自己保证自己的逻辑成功处理了,成功发了消息,但是优惠券系统,积分系统等等这么多系统,他们成功还是失败不能不考虑到

所有的服务都成功才能算这一次下单是成功的,那怎么才能保证数据一致性呢?分布式事务: 把下单,优惠券,积分。。。都放在一个事务里面一样,要成功一起成功,要失败一起失败


消息可靠性传输

消息可靠性通常要拆成三个问题来看:消息有没有成功发到 Broker、消息有没有在 Broker 中可靠存住、消息有没有被消费者正确处理完。三者任何一个环节出问题,都可能表现成“消息丢了”。

sequenceDiagram
    participant P as Producer
    participant B as Broker
    participant C as Consumer

    P->>B: 发送消息
    B-->>P: 发送确认 ack / nack
    B->>C: 投递消息

    alt 消费成功
        C-->>B: ack
        B-->>B: 标记完成或删除消息
    else 消费失败、超时或未确认
        C--x B: nack / 无 ack
        B-->>B: 重试、延迟重投或转入死信队列
    end

可以把常见治理手段概括成下面三层:

  • 生产者侧: 使用发送确认、失败重试、事务消息或本地消息表,避免“业务提交成功但消息没发出去”。
  • Broker 侧: 使用持久化、刷盘策略、副本复制和高可用部署,避免 Broker 宕机导致消息直接丢失。
  • 消费者侧: 使用手动 ack、重试、死信队列、补偿任务和幂等控制,避免“消息到了但没有被正确处理”。

从投递语义来看,大多数消息队列默认提供的是 至少一次(at least once)。这意味着重复消费通常不是异常情况,而是可靠性设计的一部分。真正的 恰好一次(exactly once) 成本很高,往往需要 Broker、存储和消费逻辑协同完成,绝大多数业务系统最终仍然要靠幂等来兜底。


JMS 和 AMQP

JMS简介:

JMS(Java Message Service,现 Jakarta Messaging)是一套面向 Java 平台的消息 API 规范。它定义的是应用程序如何以统一方式创建连接、发送消息、接收消息和处理确认,而不是底层 Broker 之间如何传输消息的网络协议。

ActiveMQ、IBM MQ 等产品都可以作为 JMS Provider 对外提供统一的 Java 消息编程模型。某些 Broker 也可以通过额外的客户端或插件提供 JMS 兼容层,但这不代表 JMS 和底层协议是同一回事。

JMS两种消息模型:

  1. 点到点(P2P)模型

JMS点到点模型.png

使用队列(Queue)作为消息通信载体;满足生产者与消费者模式,一条消息通常只会被一个消费者成功处理。 如果有多个消费者同时消费同一个队列,消息会在消费者之间分摊,而不是每个消费者都收到一份。

  1. 发布/订阅(Pub/Sub)模型

JMS发布订阅模型.png

发布订阅模型(Pub/Sub)使用主题(Topic)作为消息通信载体,类似广播模式;发布者发布一条消息后,订阅该主题的消费者都可以收到各自的一份副本。是否支持持久订阅、共享订阅以及离线后补投,要看具体规范版本和 Provider 的实现能力。

JMS 五种不同的消息正文格式:

JMS 定义了五种常见的消息体抽象,便于 Java 应用以统一 API 处理不同类型的数据:

  • StreamMessage – Java 原始值的数据流
  • MapMessage – 一组名称-值对
  • TextMessage – 字符串对象
  • ObjectMessage – 可序列化的 Java 对象
  • BytesMessage – 字节流

AMQP简介

AMQP(Advanced Message Queuing Protocol)是一种应用层消息协议,定义了客户端和消息中间件之间的线协议、消息格式和部分 Broker 行为。它关注的是消息如何在网络上传输、如何被路由、如何确认,而不是像 JMS 那样先定义一套 Java API。

RabbitMQ 是最常见的 AMQP 实现之一,尤其是 AMQP 0-9-1 模型里常见的 exchange、queue、binding、routing key、ack 等概念,很多开发者都是通过 RabbitMQ 接触到的。需要注意,AMQP 不是 JMS 的下位实现,JMS 也不是 AMQP 的上位封装,二者关注层次不同。


JMS vs AMQP

对比方向 JMS AMQP
定义 面向 Java 平台的消息 API 规范 面向客户端与 Broker 通信的应用层协议
关注点 统一编程模型 统一线协议、路由与确认语义
核心抽象 Queue、Topic、ConnectionFactory、Session Exchange、Queue、Binding、Routing Key、Ack
跨语言 主要面向 Java 天然跨语言
路由能力 主要抽象为 Queue / Topic,具体路由由 Provider 实现 协议层显式定义路由模型,可通过 direct、fanout、topic、headers 等交换器类型实现多样路由
消息体抽象 定义 TextMessage、MapMessage、BytesMessage 等消息类型 协议主要承载字节载荷和属性,业务对象通常需要自行序列化
二者关系 某些产品可对外提供 JMS 接口 某些 Broker 可在内部或外围适配 JMS,但 AMQP 本身不是 JMS

总结

  • JMS 定义的是 Java 世界里的消息编程接口,AMQP 定义的是客户端与 Broker 之间如何传输消息的协议,它们解决的是不同层面的问题
  • JMS 更强调统一的 API 抽象,AMQP 更强调跨语言通信、路由模型和确认机制
  • 在工程上可以把 JMS 看作“编程接口规范”,把 AMQP 看作“消息传输协议规范”,不要简单理解成谁替代谁

消息队列比较

消息队列比较.png


消息重复消费

消息重复消费是使用消息队列之后必须考虑的问题,也是最常见的问题之一。只要系统的目标是“尽量别丢消息”,那么重复投递和重复消费通常就不可完全避免。

从投递语义上看,大多数消息队列的可靠投递策略更接近 至少一次。也就是说,Broker 或消费者为了避免消息丢失,会选择在不确定消费结果时再次投递。因此消费者不能假设“一条消息只会被处理一次”,而要默认它可能被处理多次

就比如有这样的一个场景,用户下单成功后需要去一个活动页面给他加 GMV(销售总额),最后根据他的 GMV 去给他发奖励,这是电商活动很常见的玩法

下单了马上去看一些活动页面,有时候马上就有了,有时候缺延迟有很久,这个速度取决于消息队列的消费速度,消费慢堵塞了就迟点看到

下个单支付成功就发个消息出去,上面那个活动的开发人员就监听你的支付成功消息,监听到这个订单成功支付的消息,那就去活动 GMV 表里加上去

普通下单图解.png

但是一般消息队列的使用,都是有重试机制的,就是说下游的业务发生异常了,会抛出异常并且要求重新发一次消息

问题在于,监听消息的服务不止一个,其他服务也会有失败要求重发,但是加钱的服务是成功的,重发后,价钱的操作重复

下单异常重发.png

积分系统处理失败了,这个系统肯定要求重新发送一次这个消息,积分的系统重新接收并且处理成功了,但是别人的活动,优惠券等等服务也监听了这个消息,就可能出现活动系统给他 GMV 加两次,优惠券扣两次这种情况

真实的情况其实重试是很正常的,服务的网络抖动,开发人员代码 Bug,还有数据问题等都可能处理失败要求重发的


重复消费解决方法(接口幂等 ——— 强校验、弱校验)

幂等(idempotent、idempotence)是一个数学与计算机学概念,常见于抽象代数中。在编程中一个幂等操作的特点是其任意多次执行所产生的影响均与一次执行的影响相同

幂等函数,或幂等方法,是指可以使用相同参数重复执行,并能获得相同结果的函数。这些函数不会影响系统状态,也不用担心重复执行会对系统造成改变。例如,“setTrue()”函数就是一个幂等函数,无论多次执行,其结果都是一样的.更复杂的操作幂等保证是利用唯一交易号(流水号)实现

通俗了讲就是同样的参数调用这个接口,调用多少次结果都是一个,你加 GMV 同一个订单号加一次是多少钱,加 N 次都还是多少钱。但是如果不做幂等,一个订单调用多次钱就会出现加多次,同理退款调用多次钱也就减多次

幂等流程.png

一般幂等,会分场景考虑,是强校验还是弱校验,比如跟金钱相关的场景就很关键,就做强校验,不是很重要的场景做弱校验

强校验:

比如监听到用户支付成功的消息,监听到了去加 GMV 需要调用加钱的接口,那加钱接口下面再调用一个加流水的接口,两个放在一个事务,成功一起成功,失败一起失败

每次消息过来都要拿着 订单号 + 业务场景这样的唯一标识 (比如“天猫双十一活动加 GMV”)去流水表查,看看有没有这条流水,有就直接 return 不要走下面的流程了,没有就执行后面的逻辑

之所以用流水表,是因为涉及到金钱这样的活动,有什么问题后面也可以去流水表对账,还有就是帮助开发人员定位问题。

代码示例

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
/**
* 强校验幂等伪代码演示大致逻辑
*/
public void process(String orderId){
    try{
        // 查询订单是否存在这个活动加GMV的流水
        Object gmvFlow = getFlowByOrderId("addGmv" + orderId);
        if (Objects.isNull(gmvFlow)) {
            // 不存在流水,做加GMV和加流水操作
            addGmvAndFlow(orderId);
        } else {
            // 存在说明处理过,直接返回
            return ;
        }
    } catch (Exception e) {
        // 发送异常 触发消息队列框架重试机制
    }
}

弱校验:

一些不重要的场景,比如给谁发短信,就把这个 id + 场景唯一标识 作为 Redis 的 key,放到缓存里面失效时间看你场景,一定时间内的这个消息就去 Redis 判断

用 KV 做短时间窗口去重,适合通知、埋点这类允许少量重复或丢失也不至于造成资金、库存错误的场景。

更通用的做法还包括数据库唯一索引、去重表、状态机版本号等。Redis 更适合承担“短窗口快速去重”的职责,不适合直接承担强一致对账职责。


消息队列的顺序消费

顺序性通常讨论的是 局部有序,不是全局有序。绝大多数消息队列只能保证同一个业务键、同一个分区或同一个队列内部的顺序;一旦扩展到多个分区、多个队列或多个消费者实例,并行度提升之后,乱序问题就很容易出现。

一般是同个业务场景下不同几个操作的消息同时过去,本身顺序是对的,但是发出去的时候并发发送了,消费的时候又被并发处理,这样就会造成乱序问题

电商活动也是有这样的例子,数据量大的时候数据同步压力很大,有时候数据量大的表需要同步几个亿的数据。(并不是主从同步,主从延迟大会有问题,可能是从数据库或者主数据库同步到备库)

这种情况怼到队列里面去,然后慢慢消费,但问题是,在数据库同时对一个 Id 的数据进行了增、改、删三个操作,但是你消息发过去消费的时候变成了改,删、增,这样数据就不对,两者的结果完全不一样


顺序消费 —— RocketMQ解决案例

生产者消费者一般需要保证顺序消息的话,可能就是一个业务场景下的,比如订单的创建、支付、发货、收货

而这些东西都是同一个订单号。一个 topic 下有多个队列,为了保证发送有序,RocketMQ 提供了 MessageQueueSelector 队列选择机制,该机制有三种实现

MessageQueueSelector实现.png

可使用 Hash 取模法,让同一个订单发送到同一个队列中,再使用串行发送,只有同个订单的创建消息发送成功,再发送支付消息。这样就能保证同一个订单维度上的发送顺序

RocketMQ 能保证的是 同一消息组或同一队列内 的发送、存储和投递顺序,而不是整个 Topic 的全局顺序。只要不同订单被分散到不同队列,它们之间就可以并发消费

RocketMQ仅保证顺序发送,顺序消费由消费者业务保证

同一批需要做到顺序消费的消息根据 hash 投递到同一个 queue,这样同一个业务键就会在同一个消费实例上被顺序拉取。真正决定最终是否有序的,还是消费端是否按顺序执行和提交

但是消费者往往是多线程的,消息即使是按顺序取出的,也不代表处理一定有序。如果使用 MessageListenerOrderly,框架会帮助约束同一队列上的顺序消费;如果使用 MessageListenerConcurrently,就需要自己控制并发度,避免同一业务键被多个线程同时处理


消息队列的分布式事务

下单流程可能涉及到 10 多个环节,下单付钱都成功,但是优惠券扣减失败了,积分新增失败了,需要用分布式事务解决该问题

分布式事务大概分为:

  • 2pc(两段式提交)
  • 3pc(三段式提交)
  • TCC(Try、Confirm、Cancel)
  • 最大努力通知
  • XA
  • 本地消息表(ebay研发出的)
  • 半消息/最终一致性(RocketMQ)

这些方案并不处在完全相同的层次上。XA、2PC、3PC 更偏向协调协议;TCC 更偏向业务补偿;本地消息表、事务消息和最大努力通知更偏向最终一致性思路。

2pc(两段式):

事务两段式提交.png

2pc(两段式提交)可以理解为经典的分布式事务协调方案。它通过一个协调者先让所有参与者进入 prepare 状态,参与者先锁定资源但不提交;等协调者确认大家都准备完成之后,再统一发出 commit,否则就发出 rollback

2PC 的问题在于阻塞时间长、参与者需要长时间持有资源、协调者存在单点压力,并且在网络分区或协调者故障时容易出现状态不确定的问题

最终一致性:

事务最终一致性.png

能保证:

  • 业务主动方本地事务提交失败,业务被动方不会收到消息的投递
  • 只要业务主动方本地事务执行成功,消息就有机会被持续投递给下游;下游再通过重试、补偿、死信和人工兜底等机制,把业务状态推进到一个最终可收敛的状态

最终一致性强调的是“经过一段时间之后,系统状态能够收敛”,并不要求所有服务在同一时刻都看到完全一致的数据。工程上常见的实现包括本地消息表、事务消息、定时补偿任务和死信队列处理。