core.async 把 CSP(通信顺序进程)模型带进 Clojure——并发不靠共享内存加锁,而靠通道通信。它看起来像 Go 的 goroutine + channel,但实现与心智有微妙差异:go-block 不是真线程、无缓冲通道的交接语义、背压如何天然形成。本文从通道模型讲到死锁排查,覆盖缓冲策略、go-block 线程模型、alts! 多路选择、错误传播、超时与关闭模式、pipeline 流水线、性能陷阱,帮你把 core.async 用成「可控的并发编排工具」。
1. 通道模型:交接语义
1.1 通道是「会合点」
(require '[clojure.core.async :as async])
;; 通道三要素:放(>!)、取(<!)、关闭(close!)
(def ch (async/chan 10))
(async/put! ch :hello) ;; 非阻塞放(有缓冲)
(async/take! ch (fn [v] ...)) ;; 非阻塞取(回调)
;; go 块内用阻塞式:
(async/go (async/>! ch :x)) ;; go 块内放
(async/go (println (async/<! ch))) ;; go 块内取
心智:通道是「数据的会合点」而非「队列」——无缓冲通道下,放和取必须同时发生(handoff);有缓冲通道才是队列。理解这一点,背压就水到渠成。
1.2 无缓冲 vs 有缓冲
| 通道 | 行为 | 适用 |
|---|---|---|
(chan) | 无缓冲,放取同步交接 | 严格同步、天然背压 |
(chan 10) | 固定缓冲队列 | 削峰、解耦速率 |
(chan 10 :dropping) | 满则丢弃新值 | 丢弃可容忍的指标/日志 |
(chan 10 :sliding) | 满则丢弃旧值 | 只要「最新值」的场景 |
;; 背压的天然形成:
;; 生产者快、消费者慢 → 缓冲满 → 生产者 >! 阻塞
;; → 生产速率被消费者「拖住」(backpressure)
;; 这是特性:无界缓冲才是隐患
心法:缓冲选择 = 丢的策略。日志/指标可以
dropping(丢新的)、状态类要sliding(留新的)、业务消息要固定缓冲(满则阻塞背压,不丢不挤)。
2. go-block:不是线程
2.1 go-block 的真相
go 块不是创建线程,而是把块体编译成状态机,跑在一个共享线程池上:
;; go-block 模型
;; 1. 遇到 <! 阻塞 → 状态机挂起,让出线程
;; 2. 通道有值 → 状态机恢复,在池中任一线程继续
;; 3. 大量 go 块共享「池」线程,而非每块一线程
(async/go (async/<! ch)) ;; 挂起时不吃线程,代价低
2.2 什么时候用 thread
go 块内的代码不能阻塞在非通道操作上(sleep、数据库调用)——那会白白占用池线程:
;; 错误:阻塞式调用占住 go 池线程
(async/go (Thread/sleep 1000) ...) ;; 占线程,不好
;; 正确:真正阻塞的操作放到 thread(独立线程)
(async/thread (Thread/sleep 1000) ...)
;; 或:阻塞操作用 thread + channel 桥接
(def ch (async/to-chan (async/thread (blocking-call!))))
| 操作 | 用 go | 用 thread |
|---|---|---|
| 通道操作(<! >! alts!) | ✔ 首选 | 可 |
| 纯计算 | ✔ | — |
| 阻塞 IO / sleep / DB | ✘ | ✔ |
| 必须真实线程隔离 | — | ✔ |
心法:go-block 适合「通道驱动的非阻塞逻辑」,thread 适合「真阻塞的调用」。把 DB/IO 调用直接写进 go 块是 core.async 最常见的坑——池线程被占满,全应用卡死。
3. 背压:让慢消费拖住快生产
3.1 背压是特性
生产者 1w/s → [缓冲 10] → 消费者 1k/s
缓冲满 → 生产者 >! 阻塞 → 生产降到 1k/s(被消费者限速)
无需「流控代码」——缓冲大小天然定义「容忍的速率差」
3.2 无界缓冲的陷阱
;; 隐患:无界通道 (chan (Integer/MAX_VALUE))
;; 生产者永不阻塞 → 内存无限增长 → OOM
;; 背压 = 有界缓冲 + 阻塞放,二者缺一不可
心法:背压不是「额外的控制逻辑」,而是「有界缓冲 + 阻塞语义」自然产生的行为。生产项目里几乎永远用有界通道,让慢消费「拖住」快生产,而不是把压力转嫁为内存。
4. 多路选择:alts!
4.1 同时等多个通道
alts! 从多个操作里选先就绪的执行——超时、优先队列、扇入扇出全靠它:
;; 等多个数据源,谁先到处理谁(扇入)
(defn fan-in [sources]
(async/go-loop []
(let [[val ch] (async/alts! sources)]
(when val
(process! val)
(recur)))))
;; 带超时的等(关键模式)
(let [[val ch] (async/alts! [data-ch (async/timeout 5000)])]
(if (= ch data-ch)
(handle val)
(handle-timeout!))) ;; 超时分支
4.2 alts! 的优先级
;; 优先级:请求通道 vs 控制通道
;; 先检查控制通道(如关闭信号),再处理数据
(let [[v c] (async/alts! [ctl-ch req-ch] :priority true)]
...)
心法:
alts!是 core.async 的「select 语句」——超时(timeout通道)、多源竞争(扇入)、控制信号优先级都靠它一行解决。凡是「等多个东西」的逻辑,优先想到 alts!。
5. 错误传播:跨通道的边界
5.1 异常在通道内传播
go-block 里抛异常会传播到创建它的 go 线程(池线程),默认打日志但通道静默无值——错误不随通道走:
;; 通道内异常:不会自动传给消费者
(async/go (throw (ex-info "boom" {:x 1})))
;; 消费者永远等不到值(挂死)或拿到 nil(关闭)
;; 显式传播模式:错误「值」化
;; 约定:通道里传结果对象,含 :ok/:err 分支
(async/go (>! ch {:status :ok :data result}))
(async/go (>! ch {:status :err :error ex}))
5.2 结构化错误值
;; 统一「结果通道」约定
(defn wrap-result [f]
(fn [& args]
(async/go
(try {:status :ok :data (apply f args)}
(catch Exception e
{:status :err :error e})))))
;; 消费者侧
(async/go
(let [{:keys [status data error]} (<! ch)]
(if (= status :ok) (handle data) (log-error! error))))
心法:core.async 的通道不携带异常——错误必须显式「值化」传递。约定一个结果结构(
{:status :ok/:err}),生产消费两侧都按此解构,错误就不丢、不挂死、可日志。
6. 超时与关闭模式
6.1 timeout 通道
;; timeout 返回一个「到点自动关闭」的通道
(let [t (async/timeout 1000)]
(async/go (if (<! t) ...))) ;; 1s 后 t 关闭 → <! 返回 nil
;; 常见用途:请求超时、心跳、限时等待
6.2 优雅关闭模式
关闭通道要「约定」而非「随意」:
关闭约定:
生产者 close! → 消费者 <! 读尽缓冲后返回 nil(排干再停)
消费者不能写已关闭通道(异常)
多生产单消费:用「代理通道 + 计数」合并关闭
;; 多生产者汇合:计数关闭
(defn merge-with-count [chs]
(let [out (async/chan)]
(async/go-loop [n (count chs)]
(if (pos? n)
(let [[v _] (async/alts! chs)]
(async/>! out v)
(recur (if (nil? v) (dec n) n)))) ;; 每个源关闭 -1
(async/close! out))))
心法:超时靠
timeout通道 + alts!,关闭靠「排干 + 计数」约定——close!后<!返回 nil 是唯一的「结束信号」,用它做优雅停机(清缓冲、发完成事件、关资源)。
7. Pipeline:管道流水线
7.1 并行处理管道
;; core.async 内建 pipeline:从 in 通道读,并行处理,写 out 通道
(defn parallel-pipeline [in out n xf]
(async/pipeline n out xf in))
;; n = 并行度;xf = 转换函数(如 map)
;; 背压自动:out 满 → pipeline 停读 in
;; 手写多级流水线:采集 → 清洗 → 落库
(def stage1 (async/chan 100)) ;; 原始数据
(def stage2 (async/chan 100)) ;; 清洗后
(def stage3 (async/chan 100)) ;; 待落库
(async/pipeline 4 stage2 (map clean) stage1)
(async/pipeline 2 stage3 (map persist) stage2)
;; stage3 的消费者慢 → 逐级背压,不丢数据
7.2 并行度选择
并行度 = 每个 stage 的 worker 数
CPU 密集 → nproc(CPU 核数)
IO 密集 → 更高(等待 IO 期间让出)
下游慢 → 保持适度,靠缓冲吸收抖动
心法:pipeline 让「多级处理」显式化——每级独立通道 + 独立并行度,背压逐级传递;CPU 密集并行度取核数、IO 密集可高些。比「一个大 goroutine 里顺序做」更可控、更能榨干机器。
8. 线程安全与性能陷阱
8.1 不可变数据 + 通道 = 天然线程安全
;; 通道传的是不可变值 → 无需锁
;; 唯一可变点是通道本身(线程安全)
;; 状态更新用 atom/ref 或「通道串行化」
8.2 性能陷阱清单
| 陷阱 | 现象 | 规避 |
|---|---|---|
| go 块内阻塞调用 | 池线程占满 | 用 thread |
| 无界通道 | 内存暴涨 | 有界缓冲 |
| 大量 go 块共享池 | 饥饿 | 控制 go 数量 |
| 频繁创建通道 | GC 压力 | 复用长生命周期通道 |
| 忘记 close! | 消费者挂死 | 计数关闭约定 |
| 通道做全局总线 | 难以排查 | 按业务域分通道 |
8.3 死锁排查
死锁特征:
多个 go 块互相等对方放/取,且无超时
→ 程序「静默停住」,CPU 接近 0
排查:
1. 给每个 alts! 加 timeout(打破死锁可能)
2. 检查通道是否有界、是否 close! 了
3. 用「单生产者单消费者」先验证数据流
4. 打印通道缓冲水位(>! 前 tap 日志)
心法:死锁多半是「循环依赖 + 无超时」——给所有阻塞等待加超时(alts! + timeout),死锁就从「永久」变成「可超时」;再按「数据流分通道、逐级验证」排查。核心纪律:有界通道 + 超时 + 显式关闭三件套。
9. 实战:带超时与重试的请求编排
把前面概念串成一个真实场景——并发请求多个数据源,先到先得,整体带超时:
(defn fetch-any-with-timeout [fns timeout-ms]
(let [out (async/chan 1)
n (count fns)]
;; 每个任务一个 go:成功 → 写结果并 close,失败 → 记录继续等
(doseq [f fns]
(async/go
(try
(async/>! out (async/<! (f))) ;; 任一成功
(catch Exception _ nil))))
;; 主控制:要么拿到结果,要么超时
(async/<!!
(async/go
(let [[v _] (async/alts! [out (async/timeout timeout-ms)])]
(or v (throw (ex-info "timeout" {:ms timeout-ms}))))))))
;; 用法:三路查询,1s 内任一返回即用,超时抛错
(fetch-any-with-timeout [query-a query-b query-c] 1000)
心法:真实编排 = 扇出(doseq + go 并行查询)→ 竞争(alts! 先到先得)→ 兜底(timeout 通道)。这个模式覆盖「并发 + 超时 + 竞态收口」,是 core.async 最实用的组合拳。
10. 速查表与一句话记忆
| 问题 | 一句话答案 |
|---|---|
| 通道模型 | 交接点,无缓冲同步交接/有缓冲队列 |
| 缓冲选择 | 业务固定缓冲、指标 dropping、状态 sliding |
| go vs thread | go 非阻塞通道逻辑,thread 真阻塞调用 |
| 背压 | 有界缓冲 + 阻塞放,天然限速 |
| 多路等 | alts! + timeout 通道 |
| 错误传播 | 通道不带异常,错误值化 {:status :ok/:err} |
| 超时 | (async/timeout ms) |
| 关闭 | close! + 排干 + 计数合并 |
| 管道 | pipeline 并行处理,背压逐级传递 |
| 死锁 | 循环依赖 + 无超时 → 加超时打破 |
一句话记忆:core.async = 通道交接点(无缓冲同步/有缓冲队列)→ go-block 非阻塞状态机(真阻塞用 thread)→ 有界缓冲形成天然背压 → alts! + timeout 做多路选择与超时 → 错误必须「值化」跨通道 → 计数约定优雅关闭 → pipeline 逐级并行背压传递 → 有界 + 超时 + 显式关闭三件套防死锁——并发从「加锁防互相踩」变成「通道理顺数据流」,可控、可排查、可演进。
延伸阅读
- Clojure 并发设计模式 — STM/Agent/core.async 并发全景
- Clojure 数据管道与流处理 — core.async 在数据管道的应用
- Clojure 网络服务深入 — http-kit 异步与背压
- Clojure 微服务架构实战 — 服务间异步解耦
- Erlang 专题 — Actor 模型与 CSP 对比
- 系统架构专题 — 消息驱动与事件流架构
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。