如何提升基于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
相关产品推荐
相关产品推荐

