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

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的并发线程池资源被耗尽后,无法再消费新消息,最终表现为单条消费。

具体解决步骤

  1. 异步处理下游请求
    将同步调用下游接口的逻辑改为异步执行,释放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);
            }
        });
    }
    
  2. 启用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;
    }
    
  3. 检查事务配置
    如果消息处理方法加了@Transactional,事务未提交前ActiveMQ不会确认消息,阻塞时会导致新消息无法分发。可:

    • 移除不必要的事务,或缩小事务范围至非阻塞逻辑
    • 若必须用事务,确保事务内无长时间阻塞操作
  4. 验证监听器容器配置
    确认@JmsListener注解指定了正确的containerFactory,确保使用的是配置了setConcurrency("10-20")的工厂实例,避免默认容器覆盖配置。

内容的提问来源于stack exchange,提问作者Parth

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 11:50:32