Java多线程环境下如何仅启动一次监听器线程?含高并发优化
Java架构问题解决方案与优化建议
一、确保监听器线程仅启动一次的实现
由于没有入口main方法,只能在首次触发User.updateinfo()时初始化监听器。可以借助AtomicBoolean的原子性操作保证仅第一个调用线程启动监听器,同时将监听器设为守护线程,避免阻止JVM正常退出:
public class Data { private final ConcurrentHashMap<String, Object> dataMap = new ConcurrentHashMap<>(); private final AtomicBoolean listenerInited = new AtomicBoolean(false); private Thread monitorThread; public void add(String key, Object value) { dataMap.put(key, value); // 原子性判断并启动监听器,仅执行一次 if (listenerInited.compareAndSet(false, true)) { monitorThread = new Thread(() -> { while (!Thread.currentThread().isInterrupted()) { try { Thread.sleep(1000); // 原子性获取所有元素并清空Map,避免数据丢失 ConcurrentHashMap<String, Object> tempMap = new ConcurrentHashMap<>(); dataMap.forEach(tempMap::put); dataMap.clear(); // 打印收集到的元素 tempMap.forEach((k, v) -> System.out.printf("键:%s,值:%s%n", k, v)); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } }); monitorThread.setDaemon(true); monitorThread.start(); } } }
二、线程池优化写入操作的可行性
ConcurrentHashMap本身的put操作已做了并发优化,10000线程的并发写入它完全可以处理。如果要进一步降低锁竞争,可采用异步批量写入的思路,用单线程池做缓存合并:
public class Data { private final ConcurrentHashMap<String, Object> dataMap = new ConcurrentHashMap<>(); private final LinkedBlockingQueue<Map.Entry<String, Object>> writeQueue = new LinkedBlockingQueue<>(); private final AtomicBoolean threadsInited = new AtomicBoolean(false); private ExecutorService writeExecutor; private Thread monitorThread; public void add(String key, Object value) { writeQueue.offer(Map.entry(key, value)); if (threadsInited.compareAndSet(false, true)) { // 单线程池处理批量写入,减少Map锁竞争 writeExecutor = Executors.newSingleThreadExecutor(r -> { Thread t = new Thread(r); t.setDaemon(true); return t; }); writeExecutor.submit(() -> { while (!Thread.currentThread().isInterrupted()) { List<Map.Entry<String, Object>> batch = new ArrayList<>(100); writeQueue.drainTo(batch, 100); if (!batch.isEmpty()) { batch.forEach(entry -> dataMap.put(entry.getKey(), entry.getValue())); } try { Thread.sleep(10); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } } }); // 启动监听器线程 monitorThread = new Thread(() -> { while (!Thread.currentThread().isInterrupted()) { try { Thread.sleep(1000); ConcurrentHashMap<String, Object> tempMap = new ConcurrentHashMap<>(); dataMap.forEach(tempMap::put); dataMap.clear(); tempMap.forEach((k, v) -> System.out.printf("键:%s,值:%s%n", k, v)); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } } }); monitorThread.setDaemon(true); monitorThread.start(); } } }
注意:这种优化适合允许写入短暂延迟的场景,如果业务要求写入后必须实时可见,不建议使用。
三、架构建设性优化建议
- 拆分职责:将当前Data类拆分为三个独立模块:
DataStorage:仅负责ConcurrentHashMap的原子性读写、清空操作WriteBuffer:处理批量写入的缓存与异步提交逻辑DataMonitor:单独实现定时打印、清空的监听器逻辑
- 原子化数据获取:替换遍历+清空的非原子操作,用
AtomicReference包装Map,实现原子性的全量获取与替换:private final AtomicReference<ConcurrentHashMap<String, Object>> dataMapRef = new AtomicReference<>(new ConcurrentHashMap<>()); public ConcurrentHashMap<String, Object> takeAllAndClear() { return dataMapRef.getAndSet(new ConcurrentHashMap<>()); } - 异步打印解耦:将打印操作放入独立线程池执行,避免阻塞监听器的定时周期
- 监控告警:针对并发写入量、Map负载、队列积压量添加监控,避免出现性能瓶颈或数据堆积
- 数据结构选型:如果存在重复键场景,替换为Guava的
ConcurrentHashMultimap;如果仅需有序存储,可考虑ConcurrentSkipListMap
内容的提问来源于stack exchange,提问作者dhakshna
相关产品推荐
相关产品推荐

