Clojure 并发设计模式:STM、core.async 与 Agent 实战

深入 Clojure 并发编程三大支柱:STM 软件事务内存(ref/dosync/alter/commute)的一致性保证、core.async CSP 通道与 go-block 的异步编排、Agent 的异步队列更新模型,附生产者-消费者模式、管道流水线、背压控制与超时策略等完整实战案例。

Clojure 的并发设计哲学是「分离身份(Identity)与状态(State)」,通过不可变数据结构和专用并发原语,让多线程编程回归简单。与 Java 的锁(synchronized/ReentrantLock)和 Erlang 的 Actor 消息传递不同,Clojure 提供了三种互补的并发模型——STM(软件事务内存)、Agent 和 core.async(CSP 通道),满足从协调式状态修改到流式异步管道的各类场景。


1. 并发设计哲学:Identity vs State

1.1 不变性消除共享状态问题

传统并发编程的复杂性根源在于 共享可变状态。Clojure 通过以下设计从根本上解决问题:

策略Clojure 实现效果
不可变数据结构Persistent Vector/Map/Set状态变更创建新版本,旧版本安全共享
分离身份与状态Identity 持有 State 引用状态变更通过机制而非锁控制
专门的并发原语ref/agent/atom/ core.async不同场景选择最合适的协调模型
(def current-user (atom {:name "Alice" :role "user"}))

;; swap!:原子性更新函数
(swap! current-user assoc :role "admin")
;; => {:name "Alice", :role "admin"}(注意:返回的是新值,不改旧值)

;; 旧值仍然安全可用(不可变保证)

1.2 四种并发原语选型矩阵

原语协调模型用途典型场景
atom无锁 CAS独立状态的同步更新计数器、配置缓存
refSTM 事务多 refs 间协调一致银行账户转账、库存扣减
agent异步队列非阻塞状态更新日志写入、后台任务
core.asyncCSP 通道异步消息传递流处理、事件驱动

2. STM:软件事务内存

2.1 基础概念:ref 与 dosync

STM 将数据库事务的 ACID 特性引入内存状态管理:

(def account-a (ref 1000))
(def account-b (ref 500))

;; 事务性转账
defn transfer [from to amount]
  (dosync
    (let [from-balance @from]
      (when (< from-balance amount)
        (throw (ex-info "Insufficient funds" {:balance from-balance}))))
    (alter from - amount)
    (alter to + amount)))

(transfer account-a account-b 200)
@account-a  ; => 800
@account-b  ; => 700

2.2 alter vs commute vs ref-set

操作语义重试策略适用场景
alter严格顺序执行,可冲突检测发生写冲突时事务重试依赖顺序的操作(余额检查)
commute可交换操作,乐观并发仅在提交时验证,减少重试纯累加/累减(计数器)
ref-set直接赋值同 alter初始化、重置
(def click-counter (ref 0))

;; commute 更高效(加法可交换,重试不影响最终结果)
(dosync
  (commute click-counter inc))

;; alter 在需要读取检查结果时使用
(defn withdraw [account amount]
  (dosync
    (let [balance @account]
      (if (>= balance amount)
        (alter account - amount)
        (throw (ex-info "余额不足" {}))))))

2.3 事务嵌套与隔离性

(def inventory (ref {"apple" 10 "banana" 5}))
(def total-sales (ref 0))

(defn sell-item [item price]
  (dosync
    (let [current-stock ((ensure inventory) item)]
      ;; ensure:标记读取,防止写偏(write skew)
      (when (> current-stock 0)
        (alter inventory update item dec)
        ;; 嵌套事务自动合并到外层
        (alter total-sales + price)))))

;; 并发销售测试
(dotimes [_ 100]
  (future (sell-item "apple" 2)))

(Thread/sleep 500)
@inventory       ; => {"apple" -, "banana" 5}
@total-sales     ; => 正确总计(所有事务或全提交或全回滚)

2.4 STM 的 ACI(无 D)特性

不同于数据库的 ACID,STM 保证了:

  • Atomicity:事务内所有操作要么全成功要么全失败
  • Consistency:事务前后系统状态合法
  • Isolation:事务执行过程中不被其他事务干扰
  • ⚠️ 无 Durability:内存中的状态,重启后消失(需结合持久化)

3. Agent:异步状态代理

3.1 send 与 send-off

Agent 是 STM 的异步对应物——状态更新通过队列异步执行:

(def logger-agent (agent []))

;; send:使用线程池(适合 CPU 密集型更新)
(send logger-agent conj "用户登录: Alice")
(send logger-agent conj "订单创建: #12345")

;; send-off:创建新线程(适合 IO 密集型)
(send-off logger-agent
  (fn [logs]
    (spit "app.log" (clojure.string/join "\n" logs))
    []))

;; 等待处理完成
(await logger-agent)
@logger-agent   ; => [](已写入文件后清空)

3.2 错误处理与恢复

(def risky-agent (agent 0))

;; 触发错误
(send risky-agent (fn [_] (throw (Exception. "模拟错误"))))

;; 检查错误状态
(agent-error risky-agent)
;; => #object[java.lang.Exception]

;; 重启 agent(清空错误状态)
(restart-agent risky-agent 0 :clear-actions true)

;; 自定义错误处理
(set-error-handler! risky-agent
  (fn [agent ex]
    (println "Agent 错误:" (.getMessage ex))))

3.3 Agent 与 STM 协同

(def order-count (agent 0))
(def order-total (ref 0))

(defn process-order [amount]
  (dosync
    (alter order-total + amount))
  ;; 在事务外更新 agent(避免死锁)
  (send order-count inc))

;; Agent 更新函数内可以安全读取 STM 状态
(send order-count
  (fn [n]
    (println (format "已处理 %d 单,总金额 %d" n @order-total))
    n))

4. core.async:CSP 并发模型

4.1 通道(Channel)基础

(require '[clojure.core.async :as async :refer [<! >! chan go <!! >!! close! go-loop]])

;; 创建缓冲通道
(def c (chan 10))  ; 缓冲容量 10

;; 同步阻塞写入/读取(<!!/>>!! 为主线程阻塞)
(>!! c "hello")    ; 阻塞写入
(<!! c)            ; 阻塞读取 => "hello"

;; 异步 go 块(非阻塞,内部 park)
(go
  (>! c "来自 go block"))

(go
  (println "收到:" (<! c)))

4.2 多路复用:alts! 与超时

;; alts!:等待多个通道中任意一个就绪
(def ch1 (chan))
(def ch2 (chan))

(go
  (let [[val port] (async/alts! [ch1 ch2])]
    (println "从" port "收到:" val)))

;; 超时机制
(go
  (let [[val port] (async/alts! [ch1 (async/timeout 5000)])]
    (if (= port ch1)
      (println "收到数据:" val)
      (println "超时了!"))))

4.3 go-loop:事件循环模式

(defn worker [input-ch output-ch]
  (go-loop []
    (when-let [job (<! input-ch)]
      (let [result (process-job job)]
        (>! output-ch result))
      (recur))))

;; 优雅关闭
(defn shutdown-worker [worker-ch]
  (close! worker-ch))

5. 实战模式

5.1 生产者-消费者模式

(defn producer-consumer-system [n-consumers]
  (let [jobs (chan 100)
        results (chan)]
    ;; 启动消费者
    (dotimes [i n-consumers]
      (go-loop [processed 0]
        (if-let [job (<! jobs)]
          (do
            (Thread/sleep (rand-int 100))  ; 模拟工作
            (>! results {:consumer i :job job :processed (inc processed)})
            (recur (inc processed)))
          (println (format "消费者 #%d 处理完成,共 %d 条" i processed)))))
    {:jobs jobs :results results}))

(def system (producer-consumer-system 4))

;; 投递任务
(dotimes [i 20]
  (>!! (:jobs system) i))

;; 收集结果
(go-loop [count 0]
  (when (< count 20)
    (let [result (<! (:results system))]
      (println "完成:" result)
      (recur (inc count)))))

;; 关闭系统
(close! (:jobs system))

5.2 管道流水线(Pipeline)

(defn pipeline-stage [input-ch output-ch f parallelism]
  (async/pipeline
    parallelism      ; 并行度
    output-ch        ; 输出通道
    (map f)          ; 转换函数
    input-ch         ; 输入通道
    true             ; 关闭输出当输入关闭
    (fn [err] (println "Pipeline 错误:" err)))) ; 错误处理

;; 构建三阶段流水线
(def raw-data (chan 100))
(def parsed-data (chan 100))
(def enriched-data (chan 100))

;; 阶段1:解析
(pipeline-stage raw-data parsed-data parse-json 4)

;; 阶段2:数据增强
(pipeline-stage parsed-data enriched-data enrich-record 8)

;; 阶段3:消费
(go-loop []
  (when-let [record (<! enriched-data)]
    (save-to-db record)
    (recur)))

5.3 带背压的生产者

;; 阻塞通道(容量 0)天然实现背压
(def blocking-ch (chan))

;; 或使用 sliding/dropping 缓冲策略处理溢出
(defn create-safe-producer [buffer-size]
  (chan (async/dropping-buffer buffer-size)))  ; 满时丢弃最旧数据

;; 更精细的背压:基于通道状态调整速率
(defn adaptive-producer [output-ch]
  (go-loop [rate 100]
    (>! output-ch (generate-data))
    ;; 根据通道填充度调整速率
    (let [new-rate (if (> (count output-ch) 80) (/ rate 2) (min (* rate 1.1) 1000))]
      (Thread/sleep (long (/ 1000 rate)))
      (recur new-rate))))

6. 并发模型对比与选型

场景推荐模型理由
计数器/配置缓存atom简单 CAS,无需事务
银行转账/库存扣减ref + dosync多状态协调一致性
日志记录/后台任务agent异步非阻塞,自动队列化
流处理/事件驱动core.async自然的数据流抽象
数据管道/ETLcore.async + pipeline背压控制、错误隔离
高并发 Web 处理core.asyncGo block 轻量级并发(非线程)
定时任务/轮询core.async + timeout优雅的定时器替代

6.1 三种模型的核心差异

;; Atom:独立变量的同步更新
(def counter (atom 0))
(swap! counter inc)

;; Ref:事务内协调更新
(dosync
  (alter balance-a - amount)
  (alter balance-b + amount))

;; Agent:异步队列更新
(send logger conj log-entry)

;; Channel:数据传递而非共享
(>! channel data)   ; 发送
(<! channel)        ; 接收

7. 常见陷阱与最佳实践

7.1 STM 陷阱

陷阱说明解决方案
事务过长dosync 内执行 IO 导致长时间持有锁IO 移出事务,只用 ref 放纯计算
写偏异常读取后基于值决策,但值在提交前被改使用 ensure 标记读取依赖
无限重试commute 过多导致反复 retrycommute 仅用于可交换操作

7.2 core.async 陷阱

陷阱说明解决方案
阻塞主线程在主线程使用 <! 而非 <!!区分 go block 内(<!)和外部(<!!
通道泄漏忘记 close! 导致 goroutine 泄漏使用 go-loop + when-let 模式
缓冲溢出无缓冲通道导致死锁合理设置缓冲区或使用 dropping/sliding

7.3 Agent 陷阱

  • send 在 agent 出错后会被静默丢弃——务必配置 set-error-handler!
  • await 在主线程阻塞等待——不要在高并发流程中滥用
  • Agent 状态函数内避免持有外部锁——可能导致死锁

8. 性能基准参考

操作耗时说明
swap!~100nsCAS 原子操作
dosync (单 ref)~500nsSTM 事务开销
send~1μs入队延迟
>! / <!~50ns通道操作(无竞争时)
Thread/sleep15ms+重量级上下文切换
go block park~1μs轻量级协程切换

core.async 的 go block 基于线程池调度(固定大小线程池,默认 CPU 核心数 × 2),可以在少量线程上调度数万个并发操作,远轻于 Java 线程模型。


9. 总结

Clojure 的并发设计不是「一个框架解决所有问题」,而是通过四种专门的并发原语(atom/ref/agent/channel)和不可变数据结构,让开发者根据具体场景选择最合适的工具。STM 处理协调式状态一致性,Agent 处理异步非阻塞更新,core.async 处理流式异步管道——三者互补而非互斥。

学习路径推荐实践
入门atom 实现线程安全计数器,用 future 做简单并行计算
进阶ref + dosync 实现银行转账,用 agent 做异步日志
高级core.async 搭建生产者-消费者流水线,实现背压控制

延伸阅读可参考 Clojure 现代 Web 全栈开发 了解如何在 Web 服务中集成 core.async 处理高并发请求,以及 Clojure 调用 Java 深度实践 中使用 CompletableFuture 与 core.async 通道互操作的技巧。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「clojure」更多文章

  1. Clojure 现代 Web 全栈开发:Ring、reitit 与数据库集成
  2. Clojure spec 与测试:数据验证、生成测试与属性驱动
  3. Clojure 工具链演进:Leiningen、Clojure CLI 与 deps.edn