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

如何提升基于JMS缓存连接工厂的MQ插入线程/任务数量?

解决方案:提升MQ插入的并发线程数

一、自定义@Async专用线程池,突破异步任务并发限制

Spring默认@Async使用的SimpleAsyncTaskExecutor在实际运行中可能受环境隐含限制,导致并发线程数无法提升。自定义线程池可明确控制并发能力:

@Configuration
@EnableAsync
public class AsyncConfig implements AsyncConfigurer {

    @Override
    public Executor getAsyncExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        // 核心线程数,匹配Kafka消费线程数(10)并预留冗余
        executor.setCorePoolSize(15);
        // 最大线程数,应对峰值消息量
        executor.setMaxPoolSize(30);
        // 任务队列容量,避免消息堆积溢出
        executor.setQueueCapacity(500);
        // 线程命名前缀,方便日志排查
        executor.setThreadNamePrefix("MQ-Async-");
        // 闲置线程回收超时
        executor.setKeepAliveSeconds(60);
        // 拒绝策略:线程池满时由调用线程执行,避免丢消息
        executor.setRejectedExecutionHandler(new ThreadPoolExecutor.CallerRunsPolicy());
        executor.initialize();
        return executor;
    }

    @Override
    public AsyncUncaughtExceptionHandler getAsyncUncaughtExceptionHandler() {
        return new SimpleAsyncUncaughtExceptionHandler();
    }
}

修改异步方法,指定使用自定义线程池:

@Async("getAsyncExecutor")
public void insertIntoMQ(String kafkaMessage) {
    try{
        jmsTemplate.convertAndSend(kafkaMessage);
        log.info("Message Inserted Successfully :\n" + kafkaMessage);
    } catch (Exception e) {
        log.error("Exception Inserting Message Into MQ ", e);
    }
}

二、优化MQ连接缓存配置,匹配并发需求

当前CachingConnectionFactory仅配置了会话缓存,新增连接缓存并调整会话缓存数,确保足够的MQ连接支持高并发:

@Bean
public CachingConnectionFactory cachingConnectionFactory() {
    CachingConnectionFactory cachingConnectionFactory = new CachingConnectionFactory();
    try{
        System.setProperty("com.ibm.mq.cfg.useIBMCipherMappings", "false");
        MQConnectionFactory mqConnectionFactory = new MQConnectionFactory();
        // 原有MQ配置保持不变
        mqConnectionFactory.setHostName("ca");
        mqConnectionFactory.setPort(1414);
        mqConnectionFactory.setQueueManager("QMGR");
        mqConnectionFactory.setChannel("CHANNEL");
        mqConnectionFactory.setSSLCipherSuite("TLS_RSA_WITH_AES_128_CBC_SHA256");
        mqConnectionFactory.setTransportType(WMQConstants.WMQ_CM_CLIENT);

        UserCredentialsConnectionFactoryAdapter connectionFactoryAdapter=new UserCredentialsConnectionFactoryAdapter();
        connectionFactoryAdapter.setTargetConnectionFactory(mqConnectionFactory);
        connectionFactoryAdapter.setUsername("dmin");
        connectionFactoryAdapter.setPassword("21ff");

        cachingConnectionFactory.setTargetConnectionFactory(connectionFactoryAdapter);
        // 会话缓存数不小于异步线程池核心数
        cachingConnectionFactory.setSessionCacheSize(20);
        // 开启消费者缓存,设置连接缓存数
        cachingConnectionFactory.setCacheConsumers(true);
        cachingConnectionFactory.setConnectionCacheSize(10);

    } catch (Exception e) {
        System.out.println("Exception Inserting Message Into MQ "+ e);
    }
    return cachingConnectionFactory;
}

三、优化JmsTemplate发送配置,提升异步发送效率

开启异步发送模式,避免线程阻塞等待MQ响应:

@Bean
public JmsTemplate jmsTemplate() {
    JmsTemplate jmsTemplate = new JmsTemplate(cachingConnectionFactory());
    jmsTemplate.setDefaultDestinationName("TEST.TOPIC.QUEUE");
    // 开启异步发送
    jmsTemplate.setAsyncSend(true);
    // 设置发送超时,避免线程长时间阻塞
    jmsTemplate.setSendTimeout(5000);
    // 根据业务需求设置消息是否持久化
    jmsTemplate.setDeliveryPersistent(true);
    return jmsTemplate;
}

四、匹配Kafka消费线程数(可选)

你提到有10个Kafka消费者线程,但当前配置NUM_STREAM_THREADS_CONFIG=2,建议调整为与分区数一致(10),充分发挥消费能力:

config.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 10);

验证方式

  • 查看日志中MQ-Async-前缀的线程数量,确认并发数达到预期
  • 监控MQ消息入队速率,匹配Kafka消费速率
  • 通过JVM线程栈查看MQ相关线程的运行状态

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 13:17:05