Clojure core.async 深入:通道、背压、异常与并发编排

深入 Clojure core.async 的并发工程实践:channel 通道模型与缓冲策略(unbuffered/buffered/dropping/sliding)、go-block 与线程池模型、背压与阻塞(alts! 多路选择)、错误传播与 try/catch 在通道内的边界、timeout 与超时编排、quit 与关闭模式、pipeline 管道流水线、线程安全与性能陷阱,以及死锁与泄漏排查,帮你用 core.async 写出可控、无死锁的异步并发程序。

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 threadgo 非阻塞通道逻辑,thread 真阻塞调用
背压有界缓冲 + 阻塞放,天然限速
多路等alts! + timeout 通道
错误传播通道不带异常,错误值化 {:status :ok/:err}
超时(async/timeout ms)
关闭close! + 排干 + 计数合并
管道pipeline 并行处理,背压逐级传递
死锁循环依赖 + 无超时 → 加超时打破

一句话记忆:core.async = 通道交接点(无缓冲同步/有缓冲队列)→ go-block 非阻塞状态机(真阻塞用 thread)→ 有界缓冲形成天然背压 → alts! + timeout 做多路选择与超时 → 错误必须「值化」跨通道 → 计数约定优雅关闭 → pipeline 逐级并行背压传递 → 有界 + 超时 + 显式关闭三件套防死锁——并发从「加锁防互相踩」变成「通道理顺数据流」,可控、可排查、可演进。

延伸阅读

继续阅读

探索更多技术文章

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

全部文章 返回首页

「clojure」更多文章

  1. Clojure 函数式错误处理:Result、异常与结构化错误
  2. Clojure GraphQL API 实战:lacinia、Schema、Resolver 与权限
  3. Clojure REPL 驱动开发:nREPL、热重载与交互式工作流