You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

重试完成后消息未进入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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 04:25:06