这里主要是针对springboot中的rabbitmqtemplate来说的
基础代码以查看此专栏相关文章
消费者消息确认
rabbitmq是阅后即焚机制,rabbitmq确认消息被消费者消费后会立刻删除。
而rabbitmq是通过消费者回执来确认消费者是否成功处理消息的:消费者获取消息后,应该向rabbitmq发送ack回执,表明自己已经处理消息。
设想这样的场景
- 1)rabbitmq投递消息给消费者
- 2)消费者获取消息后,返回ack给rabbitmq
- 3)rabbitmq删除消息
- 4)消费者宕机,消息尚未处理
这样,消息就丢失了。因此消费者返回ack的时机非常重要。
springamqp则允许配置三种确认模式
- manual:手动ack,需要在业务代码结束后,调用api发送ack。
- auto:自动ack,由spring监测listener代码是否出现异常,没有异常则返回ack;抛出异常则返回nack
- none:关闭ack,mq假定消费者获取消息后会成功处理,因此消息投递后立即被删除
由此可知:
- none模式下,消息投递是不可靠的,可能丢失
- auto模式类似事务机制,出现异常时返回nack,消息回滚到mq;没有异常,返回ack
- manual:自己根据业务情况,判断什么时候该ack
一般,我们都是使用默认的auto即可。
配置方式
spring:
rabbitmq:
listener:
simple:
acknowledge-mode: 确认模式消费者失败重试机制
当消费者出现异常后,消息会不断requeue(重入队)到队列,再重新发送给消费者,然后再次异常,再次requeue,无限循环,导致mq的消息处理飙升,带来不必要的压力
本地重试
我们可以利用spring的retry机制,在消费者出现异常时利用本地重试,而不是无限制的requeue到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");
}
总结
以上为个人经验,希望能给大家一个参考,也希望大家多多支持代码网。
发表评论