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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 10:22:38