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 | 独立状态的同步更新 | 计数器、配置缓存 |
ref | STM 事务 | 多 refs 间协调一致 | 银行账户转账、库存扣减 |
agent | 异步队列 | 非阻塞状态更新 | 日志写入、后台任务 |
core.async | CSP 通道 | 异步消息传递 | 流处理、事件驱动 |
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 | 自然的数据流抽象 |
| 数据管道/ETL | core.async + pipeline | 背压控制、错误隔离 |
| 高并发 Web 处理 | core.async | Go 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 过多导致反复 retry | commute 仅用于可交换操作 |
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! | ~100ns | CAS 原子操作 |
dosync (单 ref) | ~500ns | STM 事务开销 |
send | ~1μs | 入队延迟 |
>! / <! | ~50ns | 通道操作(无竞争时) |
Thread/sleep | 15ms+ | 重量级上下文切换 |
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 通道互操作的技巧。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。