RabbitMQ 消息可靠性
1、RabbitMQ 消息丢失的可能性
消息从生产者到消费者经过三个环节:生产者、MQ、消费者,任一个环节都有可能丢失消息。
1.1 生产者消息丢失场景
- 生产者发送消息时连接 MQ 失败
- 消息到达 MQ 后未找到 Exchange
- 消息到达 MQ 的 Exchange 后,未找到合适的 Queue
- 消息到达 MQ 后,处理消息的进程发生异常
1.2 MQ 导致消息丢失
- 消息到达 MQ,保存到队列后,尚未消费就突然宕机
1.3 消费者丢失
- 消息接收后尚未处理突然宕机
- 消息接收后处理过程中抛出异常
综上,要保证 MQ 的可靠性,必须从 3 个方面入手:
- 确保生产者一定把消息发送到 MQ
- 确保 MQ 不会将消息弄丢
- 确保消费者一定要处理消息
2、如何保证生产者消息的可靠性
2.1 生产者重试机制
生产者发送消息时,出现网络故障导致与 MQ 连接中断。SpringAMQP 提供了消息发送时的重试机制,当RabbitTemplate与 MQ 连接超时后,多次重试。
在生产者对应的 yml 中配置:
spring:rabbitmq:connection-timeout:1s# 设置MQ的连接超时时间template:retry:enabled:true# 开启超时重试机制initial-interval:1000ms# 失败后的初始等待时间multiplier:2# 失败后下次的等待时长倍数,下次等待时长 = initial-interval * multipliermax-attempts:3# 最大重试次数故意写错 URL 测试,可发现总共重试了 3 次:
注意:SpringAMQP 提供的重试机制是阻塞式的,重试等待过程中当前线程被阻塞。如果对业务性能有要求,建议禁用重试机制。
2.2 生产者确认机制
一般生产者与 MQ 网络连接比较稳定,基本不用考虑第一种场景。但到达 MQ 之后可能丢失的场景包括:
- 消息到达 MQ 没有找到 Exchange
- 消息到达 MQ 找到 Exchange,但没有找到 Queue
- MQ 内部处理消息进程异常
RabbitMQ 提供生产者消息确认机制,包括Publisher Confirm和Publisher Return两种。开启确认机制后,生产者发消息给 MQ,MQ 根据处理情况返回不同回执:
- 消息发送到 MQ 但路由失败:通过 Publisher Return 返回信息,同时返回 ack 表示投递成功
- 非持久化消息发送到 MQ 且入队成功:返回 ack 表示投递成功
- 持久化消息发送到 MQ,入队成功并持久化到磁盘:返回 ack 表示投递成功
- 其他情况:返回 nack,告知投递失败
其中ack和nack属于 Publisher Confirm(ack成功,nack失败);return属于 Publisher Return。默认两者都关闭,需配置开启。
2.3 实现生产者确认
2.3.1 配置 yml 开启生产者确认
spring:rabbitmq:publisher-confirm-type:correlated# 开启publisher confirm机制,并设置confirm类型publisher-returns:true# 开启publisher return机制publisher-confirm-type三种模式:
none:关闭 confirm 机制simple:同步阻塞等待 MQ 的回执correlated:MQ 异步回调返回回执(一般使用此模式)
2.3.2 定义 ReturnCallback
每个RabbitTemplate只能配置一个 ReturnCallback,可定义配置类统一配置:
packagecom.chenwen.producer.config;importlombok.AllArgsConstructor;importlombok.extern.slf4j.Slf4j;importorg.springframework.amqp.core.ReturnedMessage;importorg.springframework.amqp.rabbit.core.RabbitTemplate;importorg.springframework.context.annotation.Configuration;importjavax.annotation.PostConstruct;@Slf4j@AllArgsConstructor@ConfigurationpublicclassReturnsCallbackConfig{privatefinalRabbitTemplaterabbitTemplate;@PostConstructpublicvoidinit(){rabbitTemplate.setReturnsCallback(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());});}}2.3.3 定义 ConfirmCallback
每个消息处理逻辑不同,需单独定义 ConfirmCallback。调用RabbitTemplate.convertAndSend时多传一个CorrelationData参数。
CorrelationData包含两个核心内容:
id:消息唯一标识,MQ 对不同的消息回执以此判断,避免混淆SettableListenableFuture:回执结果的 Future 对象
调用convertAndSend时传入CorrelationData:
MQ 的回执通过 Future 返回,可提前给 Future 添加回调:
@TestvoidtestProducerConfirmCallback()throwsInterruptedException{// 创建CorrelationDataCorrelationDatacd=newCorrelationData(UUID.randomUUID().toString());cd.getFuture().addCallback(newListenableFutureCallback<CorrelationData.Confirm>(){@OverridepublicvoidonFailure(Throwableex){log.error("消息回调失败",ex);}@OverridepublicvoidonSuccess(CorrelationData.Confirmresult){log.info("收到confirm callback回执");if(result.isAck()){log.info("消息发送成功,收到ack");}else{// 消息发送失败log.error("消息发送失败,收到nack, 原因:{}",result.getReason());}}});rabbitTemplate.convertAndSend("test.direct","chenwen","hello",cd);}测试说明:
- 路由键写错(
chenwen1):路由失败,通过 Publisher Return 返回异常信息,并返回 ACK
- 路由键正确(
chenwen):不会返回 Publisher Return 信息,只返回 ACK
注意:开启生产者确认模式较消耗 MQ 性能,一般不建议开启。分析三种场景:
- 路由失败:人为编程错误
- 交换机名称错误:编程错误
- MQ 内部故障:需要处理但概率较低,仅对消息可靠性要求极高的场景才开启,一般只需开启 Publisher Confirm 处理 nack 即可
3、MQ 消息可靠性
MQ 可靠性指消息到达 MQ 还没被消费时,MQ 因重启导致消息丢失。主要包括:
- 交换机 Exchange 持久化
- 队列 Queue 持久化
- 消息本身的持久化
3.1 Exchange 交换机持久化
Durability 参数设置持久化:Durable持久化模式,Transient临时模式。
3.2 Queues 队列持久化
队列持久化在控制台 Queues 设置 Durability:Durable持久化模式,Transient临时模式。
3.3 消息的持久化
Delivery mode 参数设为 2 即持久化。
注意:若开启消息持久化且开启生产者确认模式,需等消息持久化到磁盘才发送 ACK 回执。为减少 IO,消息并非逐条持久化,而是每隔一段时间(约 100ms)批量持久化,导致 ACK 有延迟,建议生产者确认全部采用异步方式。
3.4 LazyQueue(惰性队列)
默认情况下,生产者发消息存于内存以提高效率,但某些情况会消息堆积:
- 消费者宕机或网络故障
- 生产者生产过快,超过消费者处理能力
- 消费者处理业务发生堵塞
消息堆积导致内存占用变大,触发内存预警时,RabbitMQ 将内存消息持久化到磁盘(PageOut)。PageOut 耗时会阻塞队列进程,MQ 不再处理新消息,生产者请求被阻塞。
RabbitMQ 从 3.6.0 版本起增加 Lazy Queues(惰性队列),特性:
- 接收消息后直接存磁盘而非内存
- 消费者消费时才从磁盘读取并加载到内存(懒加载)
- 支持数百万条消息存储
3.12 版本之后,LazyQueue 已成为所有队列的默认格式。官方推荐升级 MQ 到 3.12 或所有队列设为 LazyQueue。
4、消费者的可靠性
RabbitMQ 向消费者投递消息时,可能因素导致丢失:
- 投递过程网络故障
- 消费者接收后突然宕机
- 消费者已接收但处理报错导致异常
RabbitMQ 需知道消费者处理状态,失败可再次投递。
4.1 消费者确认机制
消费者处理消息后向 RabbitMQ 发送回执,告知状态,主要有三个:
ack:处理成功,RabbitMQ 从队列删除消息nack:处理失败,RabbitMQ 重新投递reject:处理失败并拒绝,RabbitMQ 从队列删除
可用 try-catch 成功返回 ack 失败返回 nack,但 SpringAMQP 已实现,配置acknowledge-mode即可:
none:不处理,投递即 ack,消息立即删除(不建议)manual:手动模式,业务代码中调用 API 发送 ack/reject,有业务入侵但灵活auto:自动模式,SpringAMQP 用 AOP 环绕增强,正常返回 ack,失败按异常返回 nack 或 reject- 业务异常:自动返回 nack
- 消息处理或校验异常:自动返回 reject
spring:rabbitmq:listener:simple:acknowledge-mode:none# 不做处理4.1.1 测试 acknowledge-mode: none 不做处理
向test.queue发一条消息,队列当前有一条消息:
消费者监听并抛MessageConversionException。debug 断点未抛异常前刷新控制台,消息已不存在(被立即 ack 删除):
4.1.2 测试 acknowledge-mode: auto 自动处理
4.1.2.1 消费者抛出消息异常
抛MessageConversionException,异常点打断点,UI 后台消息状态为Unacked:
执行完消息数量为 0,说明消息异常直接被 reject:
4.1.2.2 消费者抛出业务异常
抛RuntimeException,断点前消息为Unacked:
异常抛出后消息回到Ready状态,确保业务异常后消息可再次投递:
5、消费者失败重试机制
5.1 消费者失败重试机制
消费者异常后消息不断 requeue 到队列重新投递,若一直失败会无限循环,导致 MQ 消息处理飙升。Spring 提供消费者重试机制:本地重试而非无限 requeue。
消费者 application.yml 配置:
spring:rabbitmq:listener:simple:retry:enabled:true# 开启消费者失败重试initial-interval:1000ms# 初始失败等待时长1秒multiplier:1# 失败等待时长倍数,下次等待时长 = multiplier * last-intervalmax-attempts:3# 最大重试次数stateless:true# true无状态;false有状态。业务含事务时改为false效果:
- 消息失败后在本地重试 3 次,不再重新入队
- 本地重试 3 次后抛出
AmqpRejectAndDontRequeueException,消息被删除(回执为 reject)
5.2 失败处理策略
失败重试 3 次后消息被删除,对可靠性要求高的场景不符合。Spring 提供失败处理策略,由MessageRecovery接口定义,三种实现:
RejectAndDontRequeueRecoverer:重试耗尽返回 reject,直接丢弃(默认)ImmediateRequeueMessageRecoverer:重试耗尽返回 nack,消息重新入队RepublishMessageRecoverer:重试耗尽将失败消息投递到指定交换机
最佳策略为RepublishMessageRecoverer,重试耗尽后投递到指定交换机,后续人工处理。示例配置:
packagecom.chenwen.consumer.config;importlombok.extern.slf4j.Slf4j;importorg.springframework.amqp.core.Binding;importorg.springframework.amqp.core.BindingBuilder;importorg.springframework.amqp.core.DirectExchange;importorg.springframework.amqp.core.Queue;importorg.springframework.amqp.rabbit.core.RabbitTemplate;importorg.springframework.amqp.rabbit.retry.MessageRecoverer;importorg.springframework.amqp.rabbit.retry.RepublishMessageRecoverer;importorg.springframework.boot.autoconfigure.condition.ConditionalOnProperty;importorg.springframework.context.annotation.Bean;importorg.springframework.context.annotation.Configuration;@Slf4j@Configuration@ConditionalOnProperty(name="spring.rabbitmq.listener.simple.retry.enabled",havingValue="true")publicclassErrorConfiguration{@BeanpublicDirectExchangeerrorExchange(){returnnewDirectExchange("error.direct");}@BeanpublicQueueerrorQueue(){returnnewQueue("error.queue");}@BeanpublicBindingerrorBinding(QueueerrorQueue,DirectExchangeerrorExchange){returnBindingBuilder.bind(errorQueue).to(errorExchange).with("error");}@BeanpublicMessageRecoverermessageRecoverer(RabbitTemplaterabbitTemplate){log.debug("加载RepublishMessageRecoverer");returnnewRepublishMessageRecoverer(rabbitTemplate,"error.direct","error");}}重试 3 次耗尽后,消息放入error.queue队列:
重试次数耗尽后,MQ 信息放在
error.queue队列中,此时error.queue多了一条数据,后续人为处理或单独监听处理。