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

HashMap与ScheduledExecutorService线程安全问题及优化咨询

问题根源分析

你遇到的定时任务停止、HashMap操作异常的核心原因有两个:

  1. HashMap非线程安全:多线程并发执行applyChanges中的put/containsKey,以及定时任务中的values遍历/clear时,会触发HashMap内部数据结构损坏(比如哈希表扩容冲突、并发修改异常),甚至导致死循环或未捕获的RuntimeException。
  2. ScheduledExecutorService的调度规则:scheduleAtFixedRate的任务一旦抛出未捕获异常,调度器会直接终止该任务的后续执行,这就是你看到定时任务不再运行的直接原因。

优化方案

方案1:原子容器交换(高效无锁,推荐)

核心思路是用AtomicReference持有缓存Map,定时任务执行时原子替换当前Map为新的空Map,再单独处理旧Map的内容;applyChanges通过CAS操作保证修改的原子性,彻底避免并发冲突。

import java.time.Duration;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;

public class StatisticEventsDispatcher {
    // 原子持有当前缓存Map
    private final AtomicReference<Map<String, AbstractStatisticsEvent>> mappedCachedEvents = 
        new AtomicReference<>(new HashMap<>());

    final Duration timeout = Duration.ofMinutes(1);
    final ScheduledExecutorService executor = Executors.newScheduledThreadPool(10);

    public StatisticEventsDispatcher(EventBus eventBus) {
        executor.scheduleAtFixedRate(
            () -> {
                try {
                    // 原子替换当前Map为新空Map,获取旧Map
                    Map<String, AbstractStatisticsEvent> oldMap = mappedCachedEvents.getAndSet(new HashMap<>());
                    // 批量发送旧Map中的事件
                    oldMap.values().forEach(eventBus::post);
                } catch (Exception e) {
                    // 捕获所有异常,避免任务终止
                    e.printStackTrace();
                }
            },
            timeout.toMillis(),
            timeout.toMillis(),
            TimeUnit.MILLISECONDS);
    }

    public void applyChanges(String type, Map<String, Long> changes) {
        while (true) {
            Map<String, AbstractStatisticsEvent> currentMap = mappedCachedEvents.get();
            // 复制当前Map,避免修改原容器
            Map<String, AbstractStatisticsEvent> newMap = new HashMap<>(currentMap);
            
            AbstractStatisticsEvent event = newMap.get(type);
            if (event == null) {
                event = new AbstractStatisticsEvent(type);
                newMap.put(type, event);
            }
            event.apply(changes);
            
            // CAS原子替换,只有当前Map未被其他线程修改时才成功
            if (mappedCachedEvents.compareAndSet(currentMap, newMap)) {
                break;
            }
            // 替换失败则重试,直到成功
        }
    }
}

方案2:线程安全容器+显式锁(简单易理解)

用ConcurrentHashMap作为线程安全容器,配合ReentrantLock锁住所有对Map的操作,保证同一时间只有一个线程能修改或遍历Map,彻底避免并发冲突。

import java.time.Duration;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.ReentrantLock;

public class StatisticEventsDispatcher {
    private final Map<String, AbstractStatisticsEvent> mappedCachedEvents = new ConcurrentHashMap<>();
    private final ReentrantLock lock = new ReentrantLock();

    final Duration timeout = Duration.ofMinutes(1);
    final ScheduledExecutorService executor = Executors.newScheduledThreadPool(10);

    public StatisticEventsDispatcher(EventBus eventBus) {
        executor.scheduleAtFixedRate(
            () -> {
                lock.lock();
                try {
                    // 批量发送事件后清空Map
                    mappedCachedEvents.values().forEach(eventBus::post);
                    mappedCachedEvents.clear();
                } catch (Exception e) {
                    e.printStackTrace();
                } finally {
                    lock.unlock();
                }
            },
            timeout.toMillis(),
            timeout.toMillis(),
            TimeUnit.MILLISECONDS);
    }

    public void applyChanges(String type, Map<String, Long> changes) {
        lock.lock();
        try {
            AbstractStatisticsEvent event = mappedCachedEvents.get(type);
            if (event == null) {
                event = new AbstractStatisticsEvent(type);
                mappedCachedEvents.put(type, event);
            }
            event.apply(changes);
        } finally {
            lock.unlock();
        }
    }
}

额外注意事项
  1. 异常捕获:无论用哪种方案,都必须给定时任务的run方法加上全局异常捕获,否则eventBus.post()抛出的异常仍会导致任务终止。
  2. Event内部线程安全:AbstractStatisticsEvent的apply方法如果涉及内部状态修改,也要保证线程安全(比如加锁、用原子变量),否则即使Map线程安全,Event内部仍会出现并发问题。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 04:50:28