数据在现代系统里是「流动的」:业务事件从客户端产生,经过管道清洗、转换、聚合,落库或进入下游系统。Clojure 的不可变数据与 core.async 通道,让数据管道天然契合「输入 → 变换 → 输出」的流水线模型。本文以 Kafka 为中心,讲解 Clojure 数据管道的完整构建:生产消费、异步管道、背压、ETL、幂等与故障恢复。
1. 数据管道模型
1.1 管道的三个阶段
来源(Source)→ 变换(Transform)→ 汇(Sink)
Kafka 主题 清洗/聚合/富化 落库/下游系统
每个阶段都是一个「函数 + 通道」的组合,阶段间用异步通道解耦:
;; 管道骨架:source → 变换函数 → sink
(defn pipeline [source-fn transform-fn sink-fn]
(let [in (async/chan 100)
out (async/chan 100)]
(async/go-loop []
(when-let [x (async/<! in)]
(async/>! out (transform-fn x))
(recur)))
(async/go-loop []
(when-let [x (async/<! out)]
(sink-fn x)
(recur)))
{:in in :out out}))
1.2 为什么用通道
通道带来解耦 + 背压:生产者不直接调用消费者,而是投递到有界缓冲区;缓冲区满时自动阻塞(背压),避免数据洪峰打垮下游。
2. Kafka 生产与消费
2.1 生产者
(require '[clj-kafka.producer :as producer])
(def producer-config
{"bootstrap.servers" "localhost:9092"
"key.serializer" "org.apache.kafka.common.serialization.StringSerializer"
"value.serializer" "org.apache.kafka.common.serialization.StringSerializer"
"acks" "all"}) ; all = 强一致
(def producer (producer/creator producer-config))
;; 发送业务事件
(producer/send producer
{:topic "order.events"
:key (str order-id)
:value (json/write-str {:order_id order-id :amount 100})})
2.2 消费者(带位移管理)
(require '[clj-kafka.consumer :as consumer])
(def consumer-config
{"bootstrap.servers" "localhost:9092"
"group.id" "order-analytics"
"key.deserializer" "org.apache.kafka.common.serialization.StringDeserializer"
"value.deserializer" "org.apache.kafka.common.serialization.StringDeserializer"
"enable.auto.commit" "false"}) ; 手动提交,配合幂等处理
(defn consume! []
(consumer/with-consumer [consumer consumer-config]
(consumer/subscribe consumer #{:order.events})
(doseq [record (consumer/poll consumer 1000)]
(process-event! (json/parse-string (:value record) true))
;; 处理成功后提交位移,失败则不提交(会重投递)
(consumer/commit! consumer))))
消费要点:
enable.auto.commit=false+ 处理成功再 commit——这是「至少一次 + 幂等处理」的标准组合,能避免处理失败却提交位移导致的数据丢失。
3. 背压与并发消费
3.1 有界缓冲 + 限流
;; 通道作为背压缓冲:限制在途消息数
(def queue (async/chan 200)) ; 最多 200 条在途
(defn feeder [topic]
(async/go-loop []
(when-let [msg (async/<! (kafka-poll! topic))]
(async/>! queue msg) ; 队列满则阻塞(背压传导到 poll)
(recur))))
(defn worker [queue]
(async/go-loop []
(when-let [msg (async/<! queue)]
(process-msg! msg)
(recur))))
;; 启动 N 个 worker 并行消费,总数由通道缓冲限制
(dotimes [i 8] (worker queue))
3.2 消费组与分区
Kafka 消费组内自动分区分工,Clojure 消费者可配合分区键保证同 key 顺序处理:
topic: order.events(3 分区)
消费组 order-analytics(3 个消费者实例)
分区 0 → 消费者 A 分区 1 → 消费者 B 分区 2 → 消费者 C
key = order_id → 同一订单的事件总是进同一分区 → 顺序有保证
4. ETL:清洗、转换与加载
4.1 清洗转换
;; 标准 ETL 变换函数:净化 → 规范化 → 富化
(defn clean-record [raw]
(-> raw
(select-keys [:order_id :amount :paid_at :user_id])
(update :amount (fn [v] (if (string? v) (parse-double v) v)))
(update :paid_at #(when % (str->instant %)))
(assoc :received_at (System/currentTimeMillis))))
(defn enrich-record [rec]
(assoc rec :amount_tier (tier-of (:amount rec))
:region (lookup-region (:user_id rec))))
;; 组装 ETL pipeline
(defn etl-pipeline [source-chan sink-fn]
(->> source-chan
(map clean-record)
(filter valid?)
(map enrich-record)
(run! sink-fn)))
要点:ETL 的每个环节都是纯函数,方便用属性测试(/clojure-property-testing/)验证「任何输入 → 清洗 → 都是合法输出」。
4.2 加载(幂等落库)
;; 落库用「upsert」保证幂等:重放也不重复
(defn sink-to-db [db rec]
(jdbc/execute! db
["insert into order_analytics (order_id, amount, amount_tier, region)
values (?, ?, ?, ?)
on conflict (order_id) do update
set amount = excluded.amount,
amount_tier = excluded.amount_tier"
(:order_id rec) (:amount rec) (:amount_tier rec) (:region rec)]))
5. 批处理与增量处理
5.1 两种处理模式
| 模式 | 触发 | 延迟 | 适用 |
|---|---|---|---|
| 流式 | 事件即时 | 秒级 | 实时看板、风控 |
| 批处理 | 定时窗口 | 小时级 | 报表、重算、回填 |
| 微批次 | 固定窗口聚合 | 分钟级 | 平衡实时与成本 |
5.2 窗口聚合
;; 按分钟窗口聚合指标(用 partition-by + reduce)
(defn windowed-aggregate [events window-ms]
(->> events
(map (juxt :ts identity))
(partition-by (fn [[ts _]] (/ ts window-ms)))
(map (fn [batch]
{:window (first batch)
:total (reduce + (map (comp :amount second) batch))
:count (count batch)}))))
6. Exactly-Once 与幂等消费
6.1 语义层次
| 语义 | 含义 | 代价 |
|---|---|---|
| At-most-once | 可能丢消息 | 最低 |
| At-least-once | 不丢但可能重复 | 中(配合幂等) |
| Exactly-once | 不丢不重 | 高(事务/Kafka Streams) |
工程建议:绝大多数场景选「At-least-once + 幂等处理」——处理函数保证「同样输入重放后状态不变」,就能达到 Exactly-once 的效果而无需昂贵的分布式事务。
6.2 幂等处理模式
;; 用去重表保证幂等:记录已处理的 event_id
(defn process-idempotent [db event]
(if (processed? db (:event_id event))
:already-processed ; 幂等跳过
(do (jdbc/transaction db
(sink-to-db db event)
(mark-processed! db (:event_id event)))
:processed)))
7. 管道可观测性与恢复
7.1 观测指标
| 指标 | 含义 |
|---|---|
| 消费延迟 | 当前 offset 与最新 offset 的差 |
| 处理速率 | 每秒处理消息数 |
| 失败率 | 处理失败与重试 |
| 通道背压 | queue 当前积压深度 |
;; 暴露管道指标(Prometheus)
(prom/gauge :pipeline_backlog "管道积压" (fn [] (async/count queue)))
7.2 故障恢复
- 死信队列:处理失败超过 N 次 → 转入 DLQ 主题,人工排查
- 重试 + 退避:瞬时失败(下游 5xx)指数退避重试
- 位移重置:代码 bug 导致消费偏移错误时,用 CLI 重置 offset 重放
- 幂等兜底:任何重放都安全(upsert + event_id 去重)
8. 常见陷阱
| 陷阱 | 现象 | 规避 |
|---|---|---|
| 自动提交位移 | 处理失败丢数据 | 手动 commit 或幂等 |
| 无界缓冲 | 内存溢出 | 通道设上限 |
| 同 key 乱序 | 状态错乱 | 按 key 分区保证顺序 |
| 反序列化失败 | 一条坏消息卡死消费 | try/catch + DLQ |
| 重放不幂等 | 数据翻倍 | 全链路幂等设计 |
9. 总结
Clojure 数据管道 = Kafka 做可靠传输 + core.async 做异步管道与背压 + 纯函数做 ETL 变换 + 幂等 upsert 做安全落库。工程红线:手动管理位移、有界缓冲传导背压、同 key 保序、坏消息进 DLQ、任何重放都幂等。管道是「数据可靠流动」的基础设施,稳住了它,/clojure-data-science/ 的分析建模才有可信的数据源。
延伸阅读
- Clojure 数据科学与数值计算 — 下游的分析与建模
- Clojure 并发设计模式 — core.async 通道与背压
- Clojure 网络服务深入 — 管道各节点的通信能力
- Kafka 专题 — Kafka 原理与可靠性
- 流处理测试 — 流处理语义(testing 专题)
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。