Spring AMQP异常时消息未转发至DLQ且保持unacked状态如何解决
当前配置使用AcknowledgeMode.MANUAL手动消息确认模式,监听器代码仅在业务逻辑完全执行成功时才会调用channel.basicAck()提交确认。一旦messageUpdater.process(message)抛出未捕获的异常,后续的ack逻辑不会执行,且代码中未定义异常场景下的nack处理逻辑,消息会一直保持unacked状态被当前消费者持有,既不会重入队列也不会路由到死信队列,只有当消费者连接断开(比如应用停止)时,消息才会被Broker重新投递给其他消费者。
RabbitMQ死信路由的触发前提是满足以下任一条件:消息被消费者明确reject/nack且设置requeue=false、消息TTL过期、队列达到最大长度。当前场景中异常发生时从未向Broker发送nack指令,自然无法触发死信投递逻辑。
根据业务需求二选一即可:
方案1:保留手动确认模式,补充异常分支的nack逻辑
如果需要手动控制确认时机,直接在监听器中增加异常捕获逻辑,业务执行失败时主动调用nack方法,设置requeue=false,Broker收到该指令后就会自动将符合死信规则的消息路由到绑定的DLQ:
@RabbitListener( queues = {"${rabbit.updater.consuming.queue.name}"}, containerFactory = "rabbitListenerContainerFactory" ) @Override public void listen( @Valid @Payload MessageDTO message, Channel channel, @Header(AmqpHeaders.DELIVERY_TAG) Long deliveryTag ) throws IOException { log.debug(DebugMessagesConstants.RECEIVED_MESSAGE_FROM_QUEUE, message, deliveryTag); try { messageUpdater.process(message); channel.basicAck(deliveryTag, false); log.debug(DebugMessagesConstants.PROCESSED_MESSAGE_FROM_QUEUE, message, deliveryTag); } catch (Exception e) { // 第三个参数为false表示不将消息重新放回原队列,直接触发死信路由 channel.basicNack(deliveryTag, false, false); log.error("消息处理失败,已投递至死信队列,消息内容:{}", message, e); } }
注意
basicNack的第二个参数为是否批量确认投递标签小于当前值的所有消息,单条消息消费场景下设为false即可。
方案2:改用自动确认模式,由Spring容器托管ack/nack逻辑(推荐)
手动确认模式需要自行覆盖所有正常、异常分支的确认逻辑,容易出现遗漏。你已经在容器工厂中配置了defaultRequeueRejected=false,只要将确认模式改为AcknowledgeMode.AUTO,SimpleRabbitListenerContainer会自动完成确认动作:业务执行成功时自动发送ack,业务抛出未捕获异常时自动发送nack且设置requeue=false,直接触发死信路由,无需手动编写channel操作逻辑。
首先修改容器工厂配置:
@Bean public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory(ObjectMapper om) { SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory(); factory.setConnectionFactory(connectionFactory()); // 改为自动确认模式,由Spring托管消息确认生命周期 factory.setAcknowledgeMode(AcknowledgeMode.AUTO); factory.setConcurrentConsumers(rabbitProperties.getUpdater().getConcurrentConsumers()); factory.setMaxConcurrentConsumers(rabbitProperties.getUpdater().getMaxConcurrentConsumers()); factory.setMessageConverter(new Jackson2JsonMessageConverter(om)); factory.setAutoStartup(rabbitProperties.getUpdater().getAutoStartup()); // 该配置指定异常场景下不重入队列,直接触发死信路由 factory.setDefaultRequeueRejected(false); return factory; }
监听器代码可以直接简化,不需要注入Channel和投递标签参数,也不需要手动调用确认方法:
@RabbitListener( queues = {"${rabbit.updater.consuming.queue.name}"}, containerFactory = "rabbitListenerContainerFactory" ) @Override public void listen(@Valid @Payload MessageDTO message) { log.debug(DebugMessagesConstants.RECEIVED_MESSAGE_FROM_QUEUE, message); messageUpdater.process(message); log.debug(DebugMessagesConstants.PROCESSED_MESSAGE_FROM_QUEUE, message); }
如果完成上述配置后消息依然无法进入死信队列,需要检查原队列的死信参数配置是否正确:
- 原业务队列必须配置
x-dead-letter-exchange参数,指定对应的死信交换机名称 - 如有自定义路由需求,需要配置
x-dead-letter-routing-key参数指定死信路由键 - 死信交换机必须和死信队列通过对应路由键完成绑定,路由规则和普通交换机一致
内容的提问来源于stack exchange,提问作者morohon

