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

消费者消息确认
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消费者确认和重试机制的核心内容。理解这些机制,能帮你在实际项目中更好地保障消息的可靠性,避免数据丢失或重复处理。希望这篇文章能给你一些启发。