商城首页欢迎来到中国正版软件门户

您的位置: 首页 > 文章列表 > 编程开发 > SpringBootRabbitMQTemplate消费者确认机制怎么选?

SpringBootRabbitMQTemplate消费者确认机制怎么选?

  发布于2026-07-24 阅读(0)

扫一扫,手机访问

这里主要针对SpringBoot中的RabbitMQTemplate来聊。

SpringBootRabbitMQTemplate消费者确认机制怎么选?

消费者消息确认

RabbitMQ的默认机制,说白了就是“阅后即焚”。一旦确认消息被消费者成功消费,它就会立刻从队列里删除。那它是怎么判断消费者是否成功处理了呢?关键就在于消费者回执——消费者拿到消息后,得向RabbitMQ发一个ACK回执,告诉它:“我处理完了,可以删了。”

来看这样一个场景

  • 1)RabbitMQ把消息投递给消费者
  • 2)消费者收到消息,返回ACK
  • 3)RabbitMQ收到ACK,删除消息
  • 4)消费者在消息还没处理完时宕机了

这会导致什么问题?消息丢了。所以,什么时候返回ACK,这个时机至关重要。

SpringAMQP提供了三种确认模式

  • manual:手动确认。业务逻辑处理完后,需要显式调用API发送ACK
  • auto:自动确认。Spring会监听消费者代码是否抛出异常,没问题就返回ACK,有异常则返回NACK
  • none:关闭确认。MQ默认消费者一定会成功处理,消息投递后就直接删了

从实际效果来看:

  • none模式最不可靠,消息可能悄无声息地丢失
  • auto模式类似于事务机制,异常就回滚,正常就提交
  • manual模式则完全由你根据业务逻辑来决定ACK的时机

大多数情况下,直接用默认的auto模式就足够了。

配置方式

spring:
  rabbitmq:
    listener:
      simple:
        acknowledge-mode: 确认模式

消费者失败重试机制

当消费者处理消息时抛出异常,消息会被重新放回队列(requeue),然后再次发送给消费者,又异常,又requeue……如此循环往复,MQ的处理压力会飙升,这显然不是什么好事。

本地重试

Spring的retry机制能很好地解决这个问题:消费者出现异常时,在本地进行重试,而不是无休止地把消息丢回MQ队列。

配置方式也很简单,在consumer服务的application.yml中添加以下内容:

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

重启consumer服务后重新测试,你会发现:

  • 重试3次后,SpringAMQP会抛出AmqpRejectAndDontRequeueException,说明本地重试生效了
  • 查看RabbitMQ控制台,消息已经被删除了——最终SpringAMQP返回的是ACK,MQ直接删除了消息

结论很清晰:

  • 开启本地重试后,消息处理异常不会requeue到队列,而是在消费者本地重试
  • 重试达到上限后,Spring会返回ACK,消息被丢弃

失败策略

从上面的测试可以看出,重试次数耗尽后,消息默认会被丢弃。这是由Spring内部的MessageRecovery接口决定的,它提供了三种不同的实现:

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

比较优雅的做法是使用RepublishMessageRecoverer,把失败的消息投递到一个专门存放异常消息的队列,后续由人工集中处理,这样既不会丢失消息,也不会影响正常流程。

具体实现分两步:

1)在consumer服务中定义处理失败消息的交换机和队列:

@Bean("error_direct")
public DirectExchange errorMessageExchange(){
    return new DirectExchange("error_direct");
}
@Bean("error_queue")
public Queue errorQueue(){
    return new Queue("error_queue", true);
}
@Bean
public Binding bindingerror(DirectExchange error_direct, Queue error_queue){
    return BindingBuilder.bind(error_queue).to(error_direct).with("error");
}

2)定义一个RepublishMessageRecoverer,关联队列和交换机:

@Bean
public MessageRecoverer messageRecoverer(RabbitTemplate rabbitTemplate){
    return new RepublishMessageRecoverer(rabbitTemplate, "error_direct", "error");
}

完整代码整合如下:

@Bean("error_direct")
public DirectExchange errorMessageExchange(){
    return new DirectExchange("error_direct");
}
@Bean("error_queue")
public Queue errorQueue(){
    return new Queue("error_queue", true);
}
@Bean
public Binding bindingerror(DirectExchange error_direct, Queue error_queue){
    return BindingBuilder.bind(error_queue).to(error_direct).with("error");
}
@Bean
public MessageRecoverer messageRecoverer(RabbitTemplate rabbitTemplate){
    return new RepublishMessageRecoverer(rabbitTemplate, "error_direct", "error");
}

总结

以上就是关于SpringBoot中RabbitMQTemplate消费者确认和重试机制的核心内容。理解这些机制,能帮你在实际项目中更好地保障消息的可靠性,避免数据丢失或重复处理。希望这篇文章能给你一些启发。

本文转载于:https://www.jb51.net/program/368018grh.htm 如有侵犯,请联系zhengruancom@outlook.com删除。
免责声明:正软商城发布此文仅为传递信息,不代表正软商城认同其观点或证实其描述。

热门关注