为何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
相关产品推荐
相关产品推荐

