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

为何Spring Boot中需用虚拟线程执行RabbitMQ异步发布?

Spring Boot中AsyncRabbitTemplate发送消息阻塞REST API的原因分析

我们的Spring Boot应用会将Controller每个请求的响应统计信息发送到RabbitMQ,当RabbitMQ服务正常时运行无问题。但我们要求这个消息发布操作是可选的:无论因任何原因导致RabbitMQ不可达或消息无法确认,都不能影响客户端请求,必须实现静默失败(仅记录错误)且不增加REST API的响应延迟。

但实际测试发现,当RabbitMQ服务器连接超时(我们在Kubernetes环境中通过禁用网络连接模拟该场景)时,AsyncRabbitTemplate::convertSendAndReceive会产生阻塞调用,直接拖慢终端用户的请求响应。

我们基于Java 21实现了一个临时解决方案:用虚拟线程Executor来异步执行消息发送操作,代码如下:

@Component
public class RabbitMQProducer {
  private static final Logger LOG = LoggerFactory.getLogger(RabbitMQProducer.class);
  private static final RateLimitedLog RATE_LIMITED_LOG = RateLimitedLog.withRateLimit(LOG).maxRate(10).every(Duration.ofSeconds(15)).build();

  private final AsyncRabbitTemplate rabbitTemplate;

  private final Executor executor = Executors.newVirtualThreadPerTaskExecutor();

  public RabbitMQProducer(AsyncRabbitTemplate rabbitTemplate) {
    this.rabbitTemplate = rabbitTemplate;
  }

  public void sendMessage(String exchangeName, String routingKey, String message) {
    executor.execute(() -> send(exchangeName, routingKey, message));
  }

  private void send(String exchangeName, String routingKey, String message) {
    final RabbitConverterFuture<String> future =
        rabbitTemplate.convertSendAndReceive(exchangeName, routingKey, message);
    future.whenComplete((result, ex) -> {
      if (ex != null) {
        RATE_LIMITED_LOG.error(message, ex);
      }
    });
  }
}

该方案在RabbitMQ无响应的场景下运行正常,但我们有一个疑问:为什么必须借助带虚拟线程的Executor?如果移除Executor,直接在sendMessage方法中调用convertSendAndReceive并监听Future,会导致阻塞进而影响用户请求?

移除Executor后的代码如下:

public void sendMessage(String exchangeName, String routingKey, String message) {
    final RabbitConverterFuture<String> future =
    rabbitTemplate.convertSendAndReceive(exchangeName, routingKey, message);
    future.whenComplete((result, ex) -> {
      if (ex != null) {
        RATE_LIMITED_LOG.error(message, ex);
      }
    });
}

以下是我们的RabbitMQ配置Bean:

@Configuration
public class RabbitMQConfig {
  @Bean
  public ConnectionFactory connectionFactory(RabbitMqConfiguration config) {
    final CachingConnectionFactory cachingConnectionFactory = new CachingConnectionFactory();

    cachingConnectionFactory.setHost(config.getAmqpHost());
    cachingConnectionFactory.setPort(config.getAmqpPort());
    cachingConnectionFactory.setUsername(config.getAmqpUser());
    cachingConnectionFactory.setPassword(config.getAmqpPassword());

    cachingConnectionFactory.getRabbitConnectionFactory().setAutomaticRecoveryEnabled(true);

    return cachingConnectionFactory;
  }

  @Bean
  public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) {
    final RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);

    rabbitTemplate.setRetryTemplate(RetryTemplate.builder().maxAttempts(1).build());

    return rabbitTemplate;
  }

  @Bean
  public AsyncRabbitTemplate asyncRabbitTemplate(RabbitTemplate rabbitTemplate) {
    return new AsyncRabbitTemplate(rabbitTemplate);
  }

  @Bean
  public SimpleRabbitListenerContainerFactory rabbitListenerContainerFactory(
      SimpleRabbitListenerContainerFactoryConfigurer configurer, ConnectionFactory connectionFactory) {
    SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
    configurer.configure(factory, connectionFactory);
    factory.setConsumerTagStrategy(q -> CONSUMER_TAG);
    return factory;
  }

  @Bean(name = "searchStreamExchange")
  public Exchange searchStreamExchange(RabbitMqConfiguration config) {
    return new TopicExchange(config.getSearchStreamExchange(), false, false);
  }
}

内容的提问来源于stack exchange,提问作者Stéphane Janel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 09:36:00