返回文章列表
高并发项目
MQRabbitMQ消息可靠性消息确认

08MQ确保消息可靠性

1.消息可靠性问题

引入 MQ 实现异步通信后,需要关注消息在传递过程中可能丢失的风险。以余额支付场景为例,支付服务扣减余额、更新支付状态后,通过 MQ 通知交易服务更新订单状态。

余额支付场景的可靠性问题

消息从发送者出发,经 MQ 中转,再到消费者处理,每个环节都可能出现问题。要保证消息可靠,需要从三方面入手:

  • 发送者的可靠性:保证发送者能把消息成功投递到 MQ。
  • MQ 的可靠性:保证 MQ 收到消息后不丢失。
  • 消费者的可靠性:保证消费者能正确处理消息并通知 MQ 删除。

2.发送者的可靠性

2.1发送者重连

有时由于网络波动,可能会出现发送者连接 MQ 失败的情况。通过配置可以开启连接失败后的重连机制:

spring:
  rabbitmq:
    connection-timeout: 1s # 设置MQ的连接超时时间
    template:
      retry:
        enabled: true # 开启超时重试机制
        initial-interval: 1000ms # 失败后的初始等待时间
        multiplier: 1 # 失败后下次的等待时长倍数,下次等待时长 = initial-interval * multiplier
        max-attempts: 3 # 最大重试次数

当网络不稳定时,利用重试机制可以有效提高消息发送的成功率。不过 SpringAMQP 提供的重试机制是阻塞式的重试,也就是说多次重试等待的过程中,当前线程是被阻塞的,会影响业务性能。如果对业务性能有要求,建议禁用重试机制;如果一定要使用,请合理配置等待时长和重试次数,也可以考虑使用异步线程来执行发送消息的代码。

2.2发送者确认

SpringAMQP 提供了 Publisher Confirm 和 Publisher Return 两种确认机制。开启确认机制后,当发送者发送消息给 MQ 后,MQ 会返回确认结果给发送者。返回的结果有以下几种情况:

  • 消息投递到了 MQ,但是路由失败。此时会通过 Publisher Return 返回路由异常原因,然后返回 ACK,告知投递成功。
  • 临时消息投递到了 MQ,并且入队成功,返回 ACK,告知投递成功。
  • 持久消息投递到了 MQ,并且入队完成持久化,返回 ACK,告知投递成功。
  • 其它情况都会返回 NACK,告知投递失败。

开启确认机制:在 publisher 微服务的 application.yml 中添加配置:

spring:
  rabbitmq:
    publisher-confirm-type: correlated # 开启publisher confirm机制,并设置confirm类型
    publisher-returns: true # 开启publisher return机制

publisher-confirm-type 有三种模式可选:

  • none:关闭 confirm 机制。
  • simple:同步阻塞等待 MQ 的回执消息。
  • correlated:MQ 异步回调方式返回回执消息。

配置 ReturnCallback:每个 RabbitTemplate 只能配置一个 ReturnCallback,因此需要在项目启动过程中配置:

@Slf4j
@AllArgsConstructor
@Configuration
public class MqConfig {
    private final RabbitTemplate rabbitTemplate;
 
    @PostConstruct
    public void init(){
        rabbitTemplate.setReturnsCallback(new RabbitTemplate.ReturnsCallback() {
            @Override
            public void returnedMessage(ReturnedMessage returned) {
                log.error("触发return callback,");
                log.debug("exchange: {}", returned.getExchange());
                log.debug("routingKey: {}", returned.getRoutingKey());
                log.debug("message: {}", returned.getMessage());
                log.debug("replyCode: {}", returned.getReplyCode());
                log.debug("replyText: {}", returned.getReplyText());
            }
        });
    }
}

发送消息时指定 ConfirmCallback:发送消息,指定消息 ID、消息 ConfirmCallback:

@Test
void testPublisherConfirm() throws InterruptedException {
    // 1.创建CorrelationData
    CorrelationData cd = new CorrelationData();
    // 2.给Future添加ConfirmCallback
    cd.getFuture().addCallback(new ListenableFutureCallback<CorrelationData.Confirm>() {
        @Override
        public void onFailure(Throwable ex) {
            // 2.1.Future发生异常时的处理逻辑,基本不会触发
            log.error("handle message ack fail", ex);
        }
        @Override
        public void onSuccess(CorrelationData.Confirm result) {
            // 2.2.Future接收到回执的处理逻辑,参数中的result就是回执内容
            if(result.isAck()){ // true代表ack回执,false 代表 nack回执
                log.debug("发送消息成功,收到 ack!");
            }else{ // result.getReason(),返回nack时的异常描述
                log.error("发送消息失败,收到 nack, reason : {}", result.getReason());
            }
        }
    });
    // 3.发送消息
    rabbitTemplate.convertAndSend("hmall.direct", "red1", "hello", cd);
}

如何处理发送者的确认消息:

  • 发送者确认需要额外的网络和系统资源开销,尽量不要使用。
  • 对于 nack 消息可以有限次数重试,依然失败则记录异常消息。

3.MQ的可靠性

在默认情况下,RabbitMQ 会将接收到的信息保存在内存中以降低消息收发的延迟。这样会导致两个问题:

  • 一旦 MQ 宕机,内存中的消息会丢失。
  • 内存空间有限,当消费者故障或处理过慢时,会导致消息积压,引发 MQ 阻塞。

3.1数据持久化

RabbitMQ 实现数据持久化包括 3 个方面:

数据持久化配置

  • 交换机持久化:MQ 重启后交换机依然存在。
  • 队列持久化:MQ 重启后队列依然存在。
  • 消息持久化:队列中的消息会持久化到磁盘,MQ 重启消息依然存在。

3.2Lazy Queue

从 RabbitMQ 的 3.6.0 版本开始,增加了 Lazy Queue(惰性队列)的概念。惰性队列的特征如下:

  • 接收到消息后直接存入磁盘,不再存储到内存。
  • 消费者要消费消息时才会从磁盘中读取并加载到内存(可以提前缓存部分消息到内存,最多 2048 条)。
  • 在 3.12 版本后,所有队列都是 Lazy Queue 模式,无法更改。

要设置一个队列为惰性队列,只需要在声明队列时指定 x-queue-mode 属性为 lazy 即可:

Lazy Queue 配置

基于代码声明 Lazy 队列:

@Bean
public Queue lazyQueue(){
    return QueueBuilder
            .durable("lazy.queue")
            .lazy() // 开启Lazy模式
            .build();
}

Lazy Queue 注解声明

基于 @RabbitListener 注解声明 Lazy 队列:

@RabbitListener(queuesToDeclare = @Queue(
    name = "lazy.queue",
    durable = "true",
    arguments = @Argument(name = "x-queue-mode", value = "lazy")
))
public void listenLazyQueue(String msg){
    log.info("接收到 lazy.queue的消息:{}", msg);
}

3.3RabbitMQ 保证消息可靠性的总结

  • 首先通过配置可以让交换机、队列、以及发送的消息都持久化。这样队列中的消息会持久化到磁盘,MQ 重启消息依然存在。
  • RabbitMQ 在 3.6 版本引入了 LazyQueue,并且在 3.12 版本后会成为队列的默认模式。LazyQueue 会将所有消息都持久化。
  • 开启持久化和发送者确认时,RabbitMQ 只有在消息持久化完成后才会给发送者返回 ACK 回执。

4.消费者的可靠性

4.1消费者确认机制

消费者确认机制(Consumer Acknowledgement)是为了确认消费者是否成功处理消息。当消费者处理消息结束后,应该向 RabbitMQ 发送一个回执,告知 RabbitMQ 自己消息处理状态:

  • ack:成功处理消息,RabbitMQ 从队列中删除该消息。
  • nack:消息处理失败,RabbitMQ 需要再次投递消息。
  • reject:消息处理失败并拒绝该消息,RabbitMQ 从队列中删除该消息。

SpringAMQP 已经实现了消息确认功能,并允许通过配置文件选择 ACK 处理方式,有三种方式:

  • none:不处理。即消息投递给消费者后立刻 ack,消息会立刻从 MQ 删除。非常不安全,不建议使用。
  • manual:手动模式。需要自己在业务代码中调用 api,发送 ack 或 reject,存在业务入侵,但更灵活。
  • auto:自动模式。SpringAMQP 利用 AOP 对消息处理逻辑做了环绕增强,当业务正常执行时则自动返回 ack。当业务出现异常时,根据异常判断返回不同结果:业务异常自动返回 nack;消息处理或校验异常自动返回 reject。
spring:
  rabbitmq:
    listener:
      simple:
        prefetch: 1
        acknowledge-mode: none # none,关闭ack;manual,手动ack;auto:自动ack

4.2失败重试机制

SpringAMQP 提供了消费者失败重试机制,在消费者出现异常时利用本地重试,而不是无限的 requeue 到 mq。通过在 application.yaml 文件中添加配置来开启重试机制:

spring:
  rabbitmq:
    listener:
      simple:
        prefetch: 1
        retry:
          enabled: true # 开启消费者失败重试
          initial-interval: 1000ms # 初始的失败等待时长为1秒
          multiplier: 1 # 下次失败的等待时长倍数,下次等待时长 = multiplier * last-interval
          max-attempts: 3 # 最大重试次数
          stateless: true # true无状态;false有状态。如果业务中包含事务,这里改为false

4.3失败消息处理策略

在开启重试模式后,重试次数耗尽,如果消息依然失败,则由 MessageRecoverer 接口来处理,它包含三种不同的实现:

  • RejectAndDontRequeueRecoverer:重试耗尽后,直接 reject,丢弃消息。默认就是这种方式。
  • ImmediateRequeueMessageRecoverer:重试耗尽后,返回 nack,消息重新入队。
  • RepublishMessageRecoverer:重试耗尽后,将失败消息投递到指定的交换机。

将失败处理策略改为 RepublishMessageRecoverer:首先定义接收失败消息的交换机、队列及其绑定关系;然后定义 RepublishMessageRecoverer:

@Bean
public MessageRecoverer republishMessageRecoverer(RabbitTemplate rabbitTemplate){
    return new RepublishMessageRecoverer(rabbitTemplate, "error.direct", "error");
}

4.4业务幂等性

幂等是一个数学概念,用函数表达式来描述是这样的:f(x) = f(f(x))。在程序开发中,则是指同一个业务,执行一次或多次对业务状态的影响是一致的。

  • 幂等:查询业务,例如根据 id 查询商品;删除业务,例如根据 id 删除商品。
  • 非幂等:用户下单业务,需要扣减库存;用户退款业务,需要恢复余额。

由于 MQ 可能会重新投递消息,消费端必须保证重复消费不会带来副作用。常用方案有两种:

方案一:唯一消息 id。给每个消息都设置一个唯一 id,利用 id 区分是否是重复消息:

  • 每一条消息都生成一个唯一的 id,与消息一起投递给消费者。
  • 消费者接收到消息后处理自己的业务,业务处理成功后将消息 ID 保存到数据库。
  • 如果下次又收到相同消息,去数据库查询判断是否存在,存在则为重复消息放弃处理。
@Bean
public MessageConverter messageConverter(){
    // 1.定义消息转换器
    Jackson2JsonMessageConverter jjmc = new Jackson2JsonMessageConverter();
    // 2.配置自动创建消息id,用于识别不同消息,也可以在业务中基于ID判断是否是重复消息
    jjmc.setCreateMessageIds(true);
    return jjmc;
}

方案二:业务判断。结合业务逻辑,基于业务本身做判断。以余额支付业务为例,处理消息的业务逻辑是把订单状态从未支付修改为已支付。因此可以在执行业务时判断订单状态是否是未支付,如果不是则证明订单已经被处理过,无需重复处理。

5.支付与交易服务的一致性保障

如何保证支付服务与交易服务之间的订单状态一致性?

  • 首先,支付服务在用户支付成功后利用 MQ 消息通知交易服务,完成订单状态同步。
  • 其次,为了保证 MQ 消息的可靠性,采用了生产者确认机制、消费者确认、消费者失败重试等策略,确保消息投递和处理的可靠性。同时也开启了 MQ 的持久化,避免因服务宕机导致消息丢失。
  • 最后,在交易服务更新订单状态时做了业务幂等判断,避免因消息重复消费导致订单状态异常。

如果交易服务消息处理失败,有没有兜底方案?可以在交易服务设置定时任务,定期查询订单支付状态。这样即便 MQ 通知失败,还可以利用定时任务作为兜底方案,确保订单支付状态的最终一致性。