DefaultJmsListenerContainerFactory并发配置不生效问题排查
问题描述
应用基于Spring JMS监听ActiveMQ队列,处理流程为:消费消息→转换格式→调用下游接口并同步等待响应→将响应发送至另一队列。配置setConcurrency("10-20")期望实现多消息并行消费,但队列积压超1000条时始终仅单条消费;注释下游请求代码后,并发消费恢复正常(10-20条并行)。
配置代码
@EnableJms @Configuration public class VmsJmsConfig implements JmsListenerConfigurer { @Value("${spring.activemq.broker-url}") String brokerUrl; @Value("${spring.activemq.ssl.trustStorePath}") String trustStorePath; @Value("${spring.activemq.ssl.trustStorePass}") String trustStorePass; @Bean public DefaultJmsListenerContainerFactory defaultJmsListenerContainerFactory(ConnectionFactory conFactory) { DefaultJmsListenerContainerFactory defaultJmsListenerContainerFactory = new DefaultJmsListenerContainerFactory(); defaultJmsListenerContainerFactory.setConnectionFactory(conFactory); defaultJmsListenerContainerFactory.setConcurrency("10-20"); return defaultJmsListenerContainerFactory; } @Bean("conFactory") public ConnectionFactory activeMQSslConnectionFactory() throws Exception { ActiveMQSslConnectionFactory activeMQSslConnectionFactory = new ActiveMQSslConnectionFactory(brokerUrl); activeMQSslConnectionFactory.setTrustStore(trustStorePath); activeMQSslConnectionFactory.setTrustStorePassword(trustStorePass); return activeMQSslConnectionFactory; } @Bean public DefaultMessageHandlerMethodFactory handlerMethodFactory() { DefaultMessageHandlerMethodFactory factory = new DefaultMessageHandlerMethodFactory(); factory.setMessageConverter(messageConverter()); return factory; } @Bean public MessageConverter messageConverter() { return new MappingJackson2MessageConverter(); } @Override public void configureJmsListeners(JmsListenerEndpointRegistrar registrar) { registrar.setMessageHandlerMethodFactory(handlerMethodFactory()); } }
排查与解决方案
核心原因
同步等待下游响应的操作阻塞了JMS监听线程,DefaultJmsListenerContainerFactory的并发线程池资源被耗尽后,无法再消费新消息,最终表现为单条消费。
具体解决步骤
异步处理下游请求
将同步调用下游接口的逻辑改为异步执行,释放JMS监听线程,使其能继续消费新消息。示例如下:// 定义自定义线程池 @Bean("downstreamThreadPool") public Executor downstreamThreadPool() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(20); executor.setMaxPoolSize(50); executor.setQueueCapacity(100); executor.setThreadNamePrefix("downstream-"); executor.initialize(); return executor; } // 在消息处理方法中异步执行阻塞逻辑 @JmsListener(destination = "your-queue-name", containerFactory = "defaultJmsListenerContainerFactory") public void handleMessage(YourMessage message) { // 非阻塞的消息转换逻辑 TransformedMsg transformedMsg = convertMessage(message); // 异步执行下游调用及响应发送 downstreamThreadPool.execute(() -> { try { DownstreamResponse response = downstreamClient.call(transformedMsg); jmsTemplate.convertAndSend("response-queue", response); } catch (Exception e) { // 异常处理:重试、死信队列投递等 handleProcessError(message, e); } }); }启用ActiveMQ连接池
默认ActiveMQSslConnectionFactory无连接池,并发线程获取连接可能阻塞,改用PooledConnectionFactory包装:@Bean("conFactory") public ConnectionFactory activeMQSslConnectionFactory() throws Exception { ActiveMQSslConnectionFactory activeMQSslConnectionFactory = new ActiveMQSslConnectionFactory(brokerUrl); activeMQSslConnectionFactory.setTrustStore(trustStorePath); activeMQSslConnectionFactory.setTrustStorePassword(trustStorePass); // 包装为连接池 PooledConnectionFactory pooledFactory = new PooledConnectionFactory(activeMQSslConnectionFactory); pooledFactory.setMaxConnections(50); // 根据并发需求调整 return pooledFactory; }检查事务配置
如果消息处理方法加了@Transactional,事务未提交前ActiveMQ不会确认消息,阻塞时会导致新消息无法分发。可:- 移除不必要的事务,或缩小事务范围至非阻塞逻辑
- 若必须用事务,确保事务内无长时间阻塞操作
验证监听器容器配置
确认@JmsListener注解指定了正确的containerFactory,确保使用的是配置了setConcurrency("10-20")的工厂实例,避免默认容器覆盖配置。
内容的提问来源于stack exchange,提问作者Parth
相关产品推荐
相关产品推荐

