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但没有找到 QueueMQ 内部处理消息进程异常RabbitMQ 提供生产者消息确认机制包括Publisher Confirm和Publisher Return两种。开启确认机制后生产者发消息给 MQMQ 根据处理情况返回不同回执消息发送到 MQ 但路由失败通过 Publisher Return 返回信息同时返回 ack 表示投递成功非持久化消息发送到 MQ 且入队成功返回 ack 表示投递成功持久化消息发送到 MQ入队成功并持久化到磁盘返回 ack 表示投递成功其他情况返回 nack告知投递失败其中ack和nack属于 Publisher Confirmack成功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 的回执correlatedMQ 异步回调返回回执一般使用此模式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;Slf4jAllArgsConstructorConfigurationpublicclassReturnsCallbackConfig{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时传入CorrelationDataMQ 的回执通过 Future 返回可提前给 Future 添加回调TestvoidtestProducerConfirmCallback()throwsInterruptedException{// 创建CorrelationDataCorrelationDatacdnewCorrelationData(UUID.randomUUID().toString());cd.getFuture().addCallback(newListenableFutureCallbackCorrelationData.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 设置 DurabilityDurable持久化模式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消息处理或校验异常自动返回 rejectspring: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说明消息异常直接被 reject4.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消息被删除回执为 reject5.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;Slf4jConfigurationConditionalOnProperty(namespring.rabbitmq.listener.simple.retry.enabled,havingValuetrue)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多了一条数据后续人为处理或单独监听处理。