
RabbitMQ 消息可靠投递confirm 确认 消费重试 死信队列导读消息队列最怕丢消息。简历投递、站内信这些场景丢一条用户就收不到反馈。我系统性地把 RabbitMQ 的可靠投递做了一遍生产端 confirm 确认、消费端手动 ack 重试、失败进死信队列。这篇记录完整的配置和踩坑。先说丢消息的三个环节生产者发送失败、MQ 自己挂了、消费者处理失败。三个环节都要兜底缺一个都不算可靠。生产端publisher confirm 确认RabbitMQ 的 confirm 机制发送消息后Broker 会异步回调确认这条消息是否成功落盘。配置spring:rabbitmq:host:127.0.0.1port:5672publisher-confirm-type:correlated# 开启 confirmpublisher-returns:true# 消息不可路由时回调template:mandatory:truepublisher-confirm-type有三个值none不开启、simple同步等待、correlated异步回调。生产用 correlated别用 simple 阻塞发消息的线程。发送消息并处理确认Slf4jComponentRequiredArgsConstructorpublicclassMqProducer{privatefinalRabbitTemplaterabbitTemplate;// 发送简历投递通知publicvoidsendResumeDelivered(ResumeDeliveredEventevent){CorrelationDatacdnewCorrelationData(UUID.randomUUID().toString());rabbitTemplate.convertAndSend(job.exchange,resume.delivered,event,cd);// 异步回调确认结果cd.getFuture().whenComplete((ack,ex)-{if(ack!nullack.isAck()){log.info(消息确认成功 id{},cd.getId());}else{log.error(消息确认失败 id{}, ack{}, cause{},cd.getId(),ack,exnull?unknown:ex.getMessage());// 落库待重发把消息存到本地表定时任务补偿resendStore.save(cd.getId(),event);}});}}关键点confirm 回调是异步的ack.isAck()为 false 表示发送失败这时消息不能直接扔掉要落一张本地表定时任务扫表重发。还有一个细节publisher-returns: truemandatory: true是处理消息发到了 exchange 但路由不到 queue的情况比如 routing key 写错。这种情况 confirm 是成功的消息确实到了 broker但永远进不了队列。需要在 RabbitTemplate 设置 ReturnsCallbackrabbitTemplate.setReturnsCallback(returned-{log.error(消息路由失败: exchange{}, routingKey{}, body{},returned.getExchange(),returned.getRoutingKey(),newString(returned.getMessage().getBody()));});消费端手动 ack 重试消费者默认是自动 ack——消息一取出来就确认不管处理成不成功。必须改成手动 ackspring:rabbitmq:listener:simple:acknowledge-mode:manual# 手动确认retry:enabled:true# 消费重试max-attempts:3# 最多重试 3 次initial-interval:1000# 重试间隔 1sdefault-requeue-rejected:falsedefault-requeue-rejected: false很关键重试 3 次还失败的消息不进原队列无限重试而是走死信。消费者代码Slf4jComponentRequiredArgsConstructorpublicclassResumeDeliveredConsumer{privatefinalResumeServiceresumeService;RabbitListener(queuesjob.resume.delivered.queue)publicvoidonMessage(ResumeDeliveredEventevent,Channelchannel,Messagemessage){longdeliveryTagmessage.getMessageProperties().getDeliveryTag();try{resumeService.handleDelivered(event);channel.basicAck(deliveryTag,false);}catch(Exceptione){log.error(处理简历投递消息失败: {},event,e);// 抛异常走 Spring 重试重试耗尽后进死信thrownewAmqpRejectAndDontRequeueException(e);}}}注意 catch 里不要自己吞异常。我之前犯过try-catch 里打了日志就返回结果消息自动 ack 了业务没处理消息丢了。正确做法处理失败抛AmqpRejectAndDontRequeueExceptionSpring 会按配置重试重试 3 次失败后拒绝入队不重新放回原队列转死信。死信队列死信队列声明ConfigurationpublicclassMqConfig{// 业务队列绑定死信交换机BeanpublicQueueresumeDeliveredQueue(){returnQueueBuilder.durable(job.resume.delivered.queue).deadLetterExchange(job.exchange.dlx).deadLetterRoutingKey(resume.delivered.dead).build();}// 死信交换机 死信队列BeanpublicDirectExchangedlxExchange(){returnnewDirectExchange(job.exchange.dlx);}BeanpublicQueueresumeDeliveredDeadQueue(){returnQueueBuilder.durable(job.resume.delivered.dead.queue).build();}BeanpublicBindingdeadBinding(){returnBindingBuilder.bind(resumeDeliveredDeadQueue()).to(dlxExchange()).with(resume.delivered.dead);}}死信队列里的消息就是重试 3 次还失败的人工排查或者定时任务捞出来补偿。踩坑记录坑一confirm 回调里直接抛异常现象ack.isAck()为 false 时我在回调里直接 throw结果消息没落库还是丢了。回调线程是 MQ 的线程池抛异常没人接业务补偿逻辑根本没执行。解决回调里只做落库和日志不抛异常。坑二ack 之后才写完数据库进程崩了现象先channel.basicAck再写业务表恰好进程崩溃消息确认了但业务没落库。解决先处理业务成功后再 ack上面代码的顺序。ack 之后业务已经完成才算真正可靠。坑三消费者异常重试 3 次全在同一秒现象重试配置了initial-interval: 1000但日志里 3 次重试间隔几乎为 0。排查发现retry.enabled配在listener.simple下但项目里用的是容器工厂手动创建的 listener配置没生效。解决确认spring.rabbitmq.listener.simple.retry的配置层级正确或者直接注入RetryInterceptorBuilder手动配置。坑四死信消息又回到原队列死循环现象死信队列消息一直堆积同时原队列也有重复消息。定位default-requeue-rejected没配重试失败的消息被重新放回原队列无限循环。解决配default-requeue-rejected: false失败进死信而不是原队列。可直接复用生产端publisher-confirm-type: correlated confirm 回调落库补偿路由失败publisher-returns: truemandatory: true ReturnsCallback消费端手动 ack业务成功才 ack失败抛AmqpRejectAndDontRequeueException重试 3 次失败进死信队列default-requeue-rejected: false死信队列消息做人工补偿或定时重放消息可靠性本质是发送确认 消费确认 兜底补偿三件套每一环都不能省。这套在求职招聘系统里扛住了线上流量照着配就行。