Spring AMQP请求回复消息:优雅关闭时延迟停止监听容器
解决方案
针对你的问题,核心是SimpleMessageListenerContainer默认在停止时,若没有正在处理的消息(队列空)会立即终止,不会等待后续可能到达的回复消息。要实现至少等待15秒的优雅关闭,可通过以下两种方式解决:
方法一:自定义容器停止逻辑(推荐)
继承SimpleMessageListenerContainer并重写doStop()方法,在正式停止前等待指定时长,期间保持队列监听以处理到达的回复:
public class WaitingReplyMessageListenerContainer extends SimpleMessageListenerContainer { // 等待时长(毫秒),对应请求消息的TTL private long waitBeforeStopMs = 0; public void setWaitBeforeStopMs(long waitBeforeStopMs) { this.waitBeforeStopMs = waitBeforeStopMs; } @Override protected void doStop() throws Exception { if (waitBeforeStopMs > 0) { long endTime = System.currentTimeMillis() + waitBeforeStopMs; // 循环等待,直到超时或线程被中断 while (System.currentTimeMillis() < endTime) { Thread.sleep(100); // 若期间收到回复并处理,容器会自然处理完成,不影响最终停止 } } // 执行原停止逻辑,关闭连接和消费者 super.doStop(); } }
然后修改你的容器配置:
@Bean public WaitingReplyMessageListenerContainer myReplyContainer( ConnectionFactory connectionFactory, RabbitTemplate myTemplate, Queue fromQueue) { WaitingReplyMessageListenerContainer container = new WaitingReplyMessageListenerContainer(); container.setConnectionFactory(connectionFactory); container.setQueues(fromQueue); container.setMessageListener(myTemplate); container.setPrefetchCount(15); container.setRecoveryInterval(15_000L); container.setDeclarationRetries(Integer.MAX_VALUE); container.setFailedDeclarationRetryInterval(15_000); // 设置等待时长为15秒,匹配请求消息TTL container.setWaitBeforeStopMs(15_000L); // 设置容器在生命周期最后停止,确保其他组件先停止发送新请求 container.setPhase(Integer.MAX_VALUE); return container; }
方法二:跟踪未完成请求并等待
维护一个未完成请求的计数器,在容器关闭时等待计数器归零或超时:
- 定义原子计数器跟踪请求状态:
@Bean public AtomicInteger pendingReplyRequests() { return new AtomicInteger(0); }
- 发送请求时递增计数器,收到回复时递减:
// 发送请求的逻辑中 CorrelationData correlationData = new CorrelationData(); pendingReplyRequests.incrementAndGet(); myTemplate.convertAndSend(exchange, routingKey, requestMessage, correlationData); // 配置RabbitTemplate的回复回调 myTemplate.setReplyCallback((correlation, reply, error) -> { pendingReplyRequests.decrementAndGet(); });
- 给回复容器添加关闭监听器,等待计数器归零或15秒超时:
@Bean public SimpleMessageListenerContainer myReplyContainer( ConnectionFactory connectionFactory, RabbitTemplate myTemplate, Queue fromQueue, AtomicInteger pendingReplyRequests) { SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(); // 原有配置... container.setPhase(Integer.MAX_VALUE); container.addShutdownListener(() -> { long startTime = System.currentTimeMillis(); // 等待未完成请求处理完毕或超时 while (pendingReplyRequests.get() > 0 && System.currentTimeMillis() - startTime < 15_000) { try { Thread.sleep(100); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } }); return container; }
关键配置补充
- 设置
container.setPhase(Integer.MAX_VALUE):确保回复容器在所有其他组件之后停止,避免新的请求被发送,同时保证已有请求的回复有足够时间被处理。 - 匹配请求消息的TTL:将等待时长设置为与请求消息的TTL一致(15秒),确保所有可能到达的回复都能被处理。
内容的提问来源于stack exchange,提问作者David Diehl
相关产品推荐
相关产品推荐

