Clojure中函数定时批量执行实现及竞态条件解决方案咨询
Clojure batch聚合函数的竞态修复方案
现有实现的竞态根因
你写的实现存在两个核心竞态来源:
- 独立维护
task和latest-batch两个atom,对两个atom的读写操作没有原子性保障:比如多个线程同时调用函数时,可能同时判断@task为nil,重复调度多个TimerTask,导致批量逻辑重复执行、参数丢失。 - TimerTask执行时,读取batch、重置状态的操作不是原子的,执行
f的过程中新进来的参数可能被写入旧batch,随后被直接清空导致丢失。
符合Clojure惯用规范的实现
Clojure处理这类需要同步的多状态场景的核心原则是:把需要原子变更的相关状态放到同一个atom里管理,依托atom的CAS原子操作语义,完全避免跨状态的竞态问题。
改进后的实现如下:
(import java.util.concurrent.Executors java.util.concurrent.TimeUnit) (defn batch-f [f delay-ms] (let [scheduler (Executors/newSingleThreadScheduledExecutor) ;; 单个atom管理所有共享状态,避免跨atom竞态 state (atom {:task nil :batch []})] (fn [& args] (let [;; 原子性把新参数加入batch,同时返回变更后的状态 new-state (swap! state update :batch conj args)] ;; 只有当前没有待执行的任务时,才调度新任务 (when-not (:task new-state) (let [new-task (fn [] ;; 原子性取出当前batch,同时重置状态 (let [batch (:batch (swap! state assoc :task nil :batch []))] (f batch))) ;; 提交调度任务 scheduled-future (.schedule scheduler new-task delay-ms TimeUnit/MILLISECONDS)] ;; 把调度后的任务实例写入state,避免重复调度 (swap! state assoc :task scheduled-future)))))))
实现说明
- 所有共享状态(待执行任务、待处理参数批次)都存在同一个atom中,所有状态变更都通过
swap!完成,天然保证原子性,不需要额外加锁 - 调度任务前通过
swap!的返回值判断是否需要创建新任务,只有第一个触发批次聚合的线程会进入调度逻辑,不会出现重复调度 - 任务执行时先原子性地取出当前批次的所有参数同时重置状态,执行
f的过程中新进来的参数会写入新的批次,不会出现参数丢失、混批的问题 - 用JDK标准的
ScheduledExecutorService替代老旧的Timer,线程管理更稳定,异常处理能力更强
内容的提问来源于stack exchange,提问作者Stepan Parunashvili
相关产品推荐
相关产品推荐

