SpringBoot多JMS监听器并发调用时日志框架线程安全问题
问题:SpringBoot多JMS监听器日志线程安全问题
背景
我们有一个SpringBoot应用,包含多个JMS监听器,从不同IBM MQ目的地读取消息。应用内置通用日志框架,可从每个监听器的输入消息中提取特定值,将每条记录以单行形式写入文件系统。
问题现象
当两个监听器同时接收消息时,日志会混入不同监听器的值。日志框架类通过@Autowired注入到两个监听器中,且每个事务开始时都会初始化值、映射和对象。
已确认JMS监听器使用不同线程,但仍出现值串用的情况,日志示例如下:
thread id ------398 name o.s.j.jlec#0-1 thread id ------401 name v2-4 04:13:17.907 [v2-4] INFO c.test.LogTransactions - springboot-test-app|946144f5|{listener="BatchPCListener", thread="o.s.j.jlec#0-1", operation="listener-1-logs"}| 04:13:17.930 [o.s.j.jlec#0-1] INFO c.test.LogTransactions - springboot-test-app|946144f5|{listener="BatchPCListener", thread="o.s.j.jlec#0-1", operation="listener-2-logs"}|
注:
o.s.j.jlec#0-1线程对应的监听器应为RealTimePCListener,但日志中错误显示为BatchPCListener(来自v2-4线程);同时该线程的operation值被更新,理论上也应出现串用问题。
已尝试的无效方案
- 为监听器分配不同的连接工厂
- 使用TaskExecutor
- 在监听器类中添加
synchronized关键字 - 用
ConcurrentMap存储日志值
需求
如何为每个监听器分配独立线程,或提供其他实现日志框架线程安全的方案?
解决方案
1. 日志框架改用ThreadLocal存储上下文
问题核心是日志框架类中的状态变量被多线程共享,即使事务初始化也会被覆盖。将状态变量改为ThreadLocal,确保每个线程持有独立副本:
@Component public class LogTransactions { private ThreadLocal<String> listenerName = new ThreadLocal<>(); private ThreadLocal<String> operation = new ThreadLocal<>(); private ThreadLocal<String> threadId = new ThreadLocal<>(); // 事务开始时初始化上下文 public void init(String listener, String op, String thread) { listenerName.set(listener); operation.set(op); threadId.set(thread); } // 生成日志 public void writeLog() { String logContent = String.format( "springboot-test-app|%s|{listener=\"%s\", thread=\"%s\", operation=\"%s\"}|", UUID.randomUUID().toString().substring(0,8), listenerName.get(), threadId.get(), operation.get() ); // 文件写入逻辑 } // 事务结束时清理,避免内存泄漏 public void clear() { listenerName.remove(); operation.remove(); threadId.remove(); } }
监听器中调用方式:
@Component public class RealTimePCListener implements MessageListener { @Autowired private LogTransactions logger; @Override @Transactional public void onMessage(Message message) { try { logger.init("RealTimePCListener", "listener-2-logs", Thread.currentThread().getName()); // 业务逻辑 logger.writeLog(); } finally { logger.clear(); } } }
2. 日志框架设为Prototype作用域
默认Spring Bean是单例,将日志框架改为原型作用域,让每个监听器持有独立实例:
@Component @Scope("prototype") public class LogTransactions { private String listenerName; private String operation; private String threadId; // 初始化、日志方法与之前一致 }
监听器中通过ObjectFactory获取实例:
@Component public class BatchPCListener implements MessageListener { @Autowired private ObjectFactory<LogTransactions> logFactory; @Override @Transactional public void onMessage(Message message) { LogTransactions logger = logFactory.getObject(); logger.init("BatchPCListener", "listener-1-logs", Thread.currentThread().getName()); // 业务逻辑 logger.writeLog(); } }
3. 重构为无状态日志框架
彻底去掉日志框架中的状态变量,将所有上下文信息作为参数直接传入日志方法:
@Component public class LogTransactions { public void writeLog(String listener, String operation, String threadId) { String logContent = String.format( "springboot-test-app|%s|{listener=\"%s\", thread=\"%s\", operation=\"%s\"}|", UUID.randomUUID().toString().substring(0,8), listener, threadId, operation ); // 文件写入逻辑 } }
监听器调用时直接传参:
@Component public class RealTimePCListener implements MessageListener { @Autowired private LogTransactions logger; @Override @Transactional public void onMessage(Message message) { // 业务逻辑处理 logger.writeLog("RealTimePCListener", "listener-2-logs", Thread.currentThread().getName()); } }
此方案完全避免线程安全问题,推荐优先采用。
4. 为每个监听器配置独立线程池
通过配置不同的JMS容器工厂,为每个监听器分配独立线程池,隔离执行环境:
@Configuration @EnableJms public class JmsConfig { @Bean public ConnectionFactory mqConnectionFactory() { // IBM MQ连接工厂配置 } // 实时监听器线程池 @Bean public TaskExecutor realTimeExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(5); executor.setMaxPoolSize(10); executor.setThreadNamePrefix("real-time-"); executor.initialize(); return executor; } // 批量监听器线程池 @Bean public TaskExecutor batchExecutor() { ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor(); executor.setCorePoolSize(5); executor.setMaxPoolSize(10); executor.setThreadNamePrefix("batch-"); executor.initialize(); return executor; } // 实时监听器容器工厂 @Bean public DefaultJmsListenerContainerFactory realTimeContainerFactory(ConnectionFactory mqConnectionFactory, TaskExecutor realTimeExecutor) { DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory(); factory.setConnectionFactory(mqConnectionFactory); factory.setTaskExecutor(realTimeExecutor); factory.setConcurrency("5-10"); return factory; } // 批量监听器容器工厂 @Bean public DefaultJmsListenerContainerFactory batchContainerFactory(ConnectionFactory mqConnectionFactory, TaskExecutor batchExecutor) { DefaultJmsListenerContainerFactory factory = new DefaultJmsListenerContainerFactory(); factory.setConnectionFactory(mqConnectionFactory); factory.setTaskExecutor(batchExecutor); factory.setConcurrency("5-10"); return factory; } }
监听器指定对应容器工厂:
@Component public class RealTimePCListener implements MessageListener { @JmsListener(destination = "REAL_TIME_QUEUE", containerFactory = "realTimeContainerFactory") @Transactional public void onMessage(Message message) { // 业务逻辑 } }
内容的提问来源于stack exchange,提问作者amiivas
相关产品推荐
相关产品推荐

