重试完成后消息未进入Dead Letter Queue(DLQ)问题求助
RabbitMQ重试耗尽后消息未进入死信队列问题
问题现象
- 消息达到配置的重试次数上限后,未被路由到Dead Letter Queue(DLQ)
- 主队列中的消息被移除,重试完成后无任何日志输出
- 重试次数符合配置,但后续无任何动作
- 尝试在消费者中抛出
AmqpRejectAndDontRequeueException,问题仍未解决
用户提供的配置与代码
队列与交换机配置
@Bean Queue testMsgDLQ() { return QueueBuilder.durable(testDLQ).build(); } @Bean DirectExchange testMsgDLX() { return new DirectExchange(testDLX); } @Bean Binding testDLBinding(Queue testMsgDLQ,DirectExchange testMsgDLX) { return BindingBuilder.bind(testMsgDLQ) .to(testMsgDLX).with(testDLQ); } @Bean Queue testMsgQueue() { return QueueBuilder.durable(testQueue) .withArgument("x-dead-letter-exchange", testDLX) .withArgument("x-dead-letter-routing-key",testDLQ).build(); } @Bean DirectExchange testMsgExchange() { return new DirectExchange(testExchange); } @Bean Binding testMessageQueue(Queue testMsgQueue,DirectExchange testMsgExchange) { return BindingBuilder.bind(testMsgQueue) .to(testMsgExchange).with(testQueue); }
容器工厂与连接配置
@Bean public ConnectionFactory connectionFactory() { CachingConnectionFactory cachingConnectionFactory = new CachingConnectionFactory(); cachingConnectionFactory.setHost(host); cachingConnectionFactory.setPort(port); cachingConnectionFactory.setUsername(username); cachingConnectionFactory.setPassword(password); cachingConnectionFactory.setAddresses(address); cachingConnectionFactory.setUri(url); cachingConnectionFactory.setCacheMode(CachingConnectionFactory.CacheMode.CONNECTION); cachingConnectionFactory.setRequestedHeartBeat(10); return cachingConnectionFactory; } @Bean public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory() { SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory(); factory.setConnectionFactory(connectionFactory()); factory.setConcurrentConsumers(concurrentConsumers); factory.setMaxConcurrentConsumers(maxConcurrentConsumers); factory.setDefaultRequeueRejected(false); factory.setAdviceChain(RetryInterceptorBuilder.stateless() .recoverer(new RejectAndDontRequeueRecoverer()) .maxAttempts(5) .backOffOptions(2000, 2.0, 10000) // initialInterval, multiplier, maxInterval .build()); return factory; } @Bean public RabbitTemplate rabbitTemplate(final ConnectionFactory publisherConnectionFactory) { final RabbitTemplate rabbitTemplate = new RabbitTemplate(publisherConnectionFactory); rabbitTemplate.setMessageConverter(producerJackson2MessageConverter()); return rabbitTemplate; } @Bean public Jackson2JsonMessageConverter producerJackson2MessageConverter() { return new Jackson2JsonMessageConverter(); } @Bean public MappingJackson2MessageConverter consumerJackson2MessageConverter() { return new MappingJackson2MessageConverter(); } @Bean public DefaultMessageHandlerMethodFactory messageHandlerMethodFactory() { DefaultMessageHandlerMethodFactory factory = new DefaultMessageHandlerMethodFactory(); factory.setMessageConverter(consumerJackson2MessageConverter()); return factory; } @Override public void configureRabbitListeners(final RabbitListenerEndpointRegistrar registrar) { registrar.setMessageHandlerMethodFactory(messageHandlerMethodFactory()); }
消费者代码
@RabbitListener(containerFactory="rabbitListenerContainerFactory", queues = "${rabbitmq.email.queue.name}") public void receiveMessage(MessageKeyVo messageId) { System.out.println("Demo"); throw new RuntimeException("q"); }
解决方案
1. 移除冲突的容器配置
删除factory.setDefaultRequeueRejected(false),因为重试拦截器的RejectAndDontRequeueRecoverer已经明确指定重试耗尽后拒绝消息并不重新入队,该配置会与恢复器的行为冲突,导致死信触发逻辑异常。
2. 验证死信路由配置一致性
- 确认
testDLX、testDLQ等变量的实际值无拼写错误 - 若队列已预先创建,需删除原有队列后重新启动应用,确保死信参数生效(RabbitMQ队列创建后参数无法修改)
- 通过RabbitMQ管理界面检查主队列的
x-dead-letter-exchange和x-dead-letter-routing-key参数是否正确,死信队列与死信交换机的绑定是否匹配
3. 改用主动重发布恢复器(推荐)
替换RejectAndDontRequeueRecoverer为RepublishMessageRecoverer,主动将重试耗尽的消息发布到死信交换机,不依赖队列的死信机制,可靠性更高:
@Bean public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory(RabbitTemplate rabbitTemplate) { SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory(); factory.setConnectionFactory(connectionFactory()); factory.setConcurrentConsumers(concurrentConsumers); factory.setMaxConcurrentConsumers(maxConcurrentConsumers); factory.setAdviceChain(RetryInterceptorBuilder.stateless() // 直接指定死信交换机和路由键 .recoverer(new RepublishMessageRecoverer(rabbitTemplate, testDLX, testDLQ)) .maxAttempts(5) .backOffOptions(2000, 2.0, 10000) .build()); return factory; }
4. 添加日志排查
在消费者和恢复器中添加日志,确认重试流程和消息状态:
@RabbitListener(containerFactory="rabbitListenerContainerFactory", queues = "${rabbitmq.email.queue.name}") public void receiveMessage(MessageKeyVo messageId, Message message) { log.info("处理消息ID: {}", message.getMessageProperties().getMessageId()); throw new RuntimeException("测试重试异常"); } // 自定义恢复器添加日志 .recoverer((message, cause) -> { log.error("消息重试耗尽,ID: {}, 原因: {}", message.getMessageProperties().getMessageId(), cause.getMessage()); new RejectAndDontRequeueRecoverer().recover(message, cause); })
内容的提问来源于stack exchange,提问作者Prithivi Thiyagu
相关产品推荐
相关产品推荐

