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

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;
}

方法二:跟踪未完成请求并等待

维护一个未完成请求的计数器,在容器关闭时等待计数器归零或超时:

  1. 定义原子计数器跟踪请求状态:
@Bean
public AtomicInteger pendingReplyRequests() {
    return new AtomicInteger(0);
}
  1. 发送请求时递增计数器,收到回复时递减:
// 发送请求的逻辑中
CorrelationData correlationData = new CorrelationData();
pendingReplyRequests.incrementAndGet();

myTemplate.convertAndSend(exchange, routingKey, requestMessage, correlationData);

// 配置RabbitTemplate的回复回调
myTemplate.setReplyCallback((correlation, reply, error) -> {
    pendingReplyRequests.decrementAndGet();
});
  1. 给回复容器添加关闭监听器,等待计数器归零或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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 18:15:20