HashMap与ScheduledExecutorService线程安全问题及优化咨询
问题根源分析
你遇到的定时任务停止、HashMap操作异常的核心原因有两个:
- HashMap非线程安全:多线程并发执行
applyChanges中的put/containsKey,以及定时任务中的values遍历/clear时,会触发HashMap内部数据结构损坏(比如哈希表扩容冲突、并发修改异常),甚至导致死循环或未捕获的RuntimeException。 - 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(); } } }
额外注意事项
- 异常捕获:无论用哪种方案,都必须给定时任务的
run方法加上全局异常捕获,否则eventBus.post()抛出的异常仍会导致任务终止。 - Event内部线程安全:
AbstractStatisticsEvent的apply方法如果涉及内部状态修改,也要保证线程安全(比如加锁、用原子变量),否则即使Map线程安全,Event内部仍会出现并发问题。
内容的提问来源于stack exchange,提问作者antpaw
相关产品推荐
相关产品推荐

