事件驱动架构里,Kafka 常是那条「主干道」,而 Clojure 服务则是路上的「处理站」。把两者接好,关键在三件事:客户端选型(Jackdaw)、序列化契约(Schema Registry)、位移与错误处理(不丢不重)。本文从 Jackdaw 客户端讲到生产运维,覆盖生产与消费、序列化、位移提交、与 core.async 的编排、与 Kafka Streams 的取舍,帮你把 Clojure 服务接进一条可靠、可观测的事件流。
1. Kafka 与 Clojure 生态
1.1 客户端选型
| 客户端 | 语言 | 特点 |
|---|---|---|
| Jackdaw | Clojure 包装 | 惯用 Clojure API,基于官方 Java 客户端 |
| 官方 Java 客户端 | Java | 直接互操作,无封装 |
| kafka-clj | 纯 Clojure | 早期项目,维护弱 |
1.2 为什么用 Jackdaw
;; Jackdaw 把 Java 客户端的 Builder 风格包成 Clojure 数据
{:bootstrap.servers "localhost:9092"
:key.serializer :string
:value.serializer :edn}
配置是 map、序列化器是关键字、返回值是 Clojure 数据——这是「用 Clojure 的方式用 Kafka」。
1.3 事件驱动的心智
核心概念映射:
topic —— 事件类别(orders、payments)
key —— 分区依据(常是实体 id,保证同实体有序)
value —— 事件负载
offset —— 分区内位置(消费进度)
group —— 消费组(分区在组内分配)
心智:Kafka 的 key 决定分区,分区决定顺序——「同一订单的事件要有序」靠的是「同一订单用同一 key」。理解这一点,乱序与重复问题就有了抓手。
2. Jackdaw 客户端
2.1 依赖
;; deps.edn
{:deps {fundingcircle/jackdaw {:mvn/version "0.9.9"}}}
2.2 创建生产者与消费者
(require '[jackdaw.client :as kc]
'[jackdaw.client.log :as klog]
'[jackdaw.admin :as ka])
(def producer-config
{:bootstrap.servers "localhost:9092"
:key.serializer :string
:value.serializer :edn})
(def consumer-config
(merge {:bootstrap.servers "localhost:9092"
:group.id "order-processor"
:key.deserializer :string
:value.deserializer :edn
:auto.offset.reset :earliest}
{}))
2.3 管理 topic
(def admin-client (ka/->AdminClient producer-config))
(ka/create-topics! admin-client
[{:topic-name "orders"
:partition-count 6
:replication-factor 3
:topic-config {"retention.ms" "604800000"}}])
心法:把 topic 当「接口」来管理——分区数、副本、保留期都是契约的一部分,用代码声明(而非手工敲命令)才能纳入版本控制与评审。
3. 生产消息
3.1 基本生产
(require '[jackdaw.client :as kc]
'[jackdaw.client.partitioning :as part])
(def producer (kc/producer producer-config))
(with-open [p producer]
(let [record {:topic-name "orders"
:key "order-1234"
:value {:order-id 1234 :amount 199.00}}]
@(kc/produce! p record)))
3.2 分区策略
;; 按 key 哈希分区(默认,保证同 key 有序)
{:topic-name "orders" :key "order-1234" :value {...}}
;; 显式指定分区
{:topic-name "orders" :partition 3 :value {...}}
;; 自定义分区器
(part/->HashPartitioner)
3.3 同步与异步
;; 异步:返回 future,不阻塞
(def f (kc/produce! producer record))
@f ;; 需要确认时再 deref
;; 批量生产:一次提交多条
(kc/produce! producer
[{:topic-name "orders" :key "o1" :value {...}}
{:topic-name "orders" :key "o2" :value {...}}])
心法:生产者的可靠性靠「acks + 重试 + 幂等」三件套——
acks=all保证写入副本、enable.idempotence=true避免重试导致的重复。默认配置对「不丢不重」并不够。
4. 消费与消费组
4.1 轮询消费
(require '[jackdaw.client.log :as klog])
(def consumer (kc/consumer consumer-config))
(kc/subscribe consumer [{:topic-name "orders"}])
(with-open [c consumer]
(loop []
(let [records (kc/poll c 1000)] ;; 1s 超时
(doseq [r records]
(println "收到" (:key r) (:value r)))
(recur))))
4.2 消费组与再均衡
再均衡(rebalance):
组内消费者增减 -> 分区重新分配
期间消费暂停(stop-the-world)
对策:增量协作式再均衡、静态成员(static membership)
;; 静态成员:重启不掉组,避免无谓再均衡
{:group.instance.id "order-processor-1"
:session.timeout.ms 30000}
4.3 提交位移
;; 自动提交(简单但可能重复/丢失)
{:enable.auto.commit true :auto.commit.interval.ms 5000}
;; 手动提交(处理完再提交,更可控)
(kc/commit-sync consumer)
心法:位移提交点是「语义边界」——先处理再提交 = at-least-once(可能重复),先提交再处理 = at-most-once(可能丢)。多数业务选 at-least-once,再用幂等兜底。
5. 序列化与 Schema
5.1 序列化器
;; 内置:string / edn / json / avro
{:value.serializer :edn}
{:value.serializer :json}
{:value.serializer :string}
5.2 与 Schema Registry 集成
;; 用 Confluent 的 Avro 序列化器接 Schema Registry
(def serde-config
{:schema.registry.url "http://schema-registry:8081"
:value.serializer :avro
:key.serializer :string})
;; 值用 Avro 编码,Schema 由 Registry 管理版本
5.3 Schema 演进
| 变更 | 兼容性 | 说明 |
|---|---|---|
| 加带默认值的字段 | 兼容 | 推荐 |
| 删字段 | 视配置 | 需默认值 |
| 改字段类型 | 不兼容 | 需新 topic |
| 改字段名 | 不兼容 | 需别名 |
心法:Schema 是生产消费双方之间的契约——用 Schema Registry 把它版本化,加字段带默认值向前兼容,改类型改名字当破坏性变更处理。没有契约的 JSON 是「今天能跑,明天崩」。
6. 位移管理与错误处理
6.1 处理与提交的原子性
(defn process-batch [consumer records]
(try
;; 1. 处理整批(业务幂等)
(doseq [r records] (handle! r))
;; 2. 处理成功后才提交
(kc/commit-sync consumer)
(catch Exception e
(log/error e "批处理失败")
(throw e))))
6.2 死信队列
(defn safe-handle! [producer r]
(try
(handle! r)
(catch Exception e
;; 处理不了的毒消息进 DLQ,不阻塞后续
@(kc/produce! producer
{:topic-name "orders-dlq"
:key (:key r)
:value {:error (ex-message e)
:original (:value r)}}))))
6.3 重试与退避
错误分类:
可重试 —— 网络抖动、依赖超时(退避重试)
不可重试 —— 数据格式错、业务规则拒绝(进 DLQ)
致命 —— 序列化失败、权限错(停止消费并告警)
心法:「毒消息」是消费端最常见的稳定性杀手——一条永远处理不了的消息会让分区卡死。把它路由到死信队列,让主干继续流动,事后人工分析。
7. 与 core.async 编排
7.1 消费到通道
(require '[clojure.core.async :as async])
(defn kafka->chan [consumer ch]
(async/thread
(try
(loop []
(let [records (kc/poll consumer 1000)]
(doseq [r records]
(async/>!! ch r)) ;; 阻塞放(有背压)
(recur)))
(finally
(async/close! ch)))))
(def events (async/chan 100))
(kafka->chan consumer events)
7.2 并行处理管道
;; 用 pipeline 并行处理,背压自动
(def processed (async/chan 100))
(async/pipeline 4 processed (map transform) events)
(async/go-loop []
(when-let [r (async/<! processed)]
(persist! r)
(recur)))
7.3 反压与提交的配合
要点:
poll 的批次 -> 有界通道(背压)
通道满 -> poll 阻塞 -> 消费自然放缓
处理完成 -> 提交位移
原则:位移提交必须「在处理真正落库之后」
心法:core.async 给 Kafka 消费加上「背压」——有界通道让下游慢时自动拖住 poll,避免「拉太快、内存爆」。这是 Clojure 消费端最优雅的流控方式。
8. 与 Kafka Streams 集成
8.1 什么时候用 Streams
| 需求 | 用 Jackdaw 手写 | 用 Kafka Streams |
|---|---|---|
| 简单消费处理 | 合适 | 过重 |
| 有状态聚合 | 要自己管状态 | 内置状态存储 |
| 窗口与 join | 手写复杂 | 原生支持 |
| exactly-once | 需自己实现 | 内置支持 |
8.2 从 Clojure 调 Streams
;; Kafka Streams 是 Java API,Clojure 通过 interop 使用
(import '[org.apache.kafka.streams StreamsBuilder]
'[org.apache.kafka.streams.kstream KStream])
(def builder (StreamsBuilder.))
(def stream (.stream builder "orders"))
(.foreach stream
(reify org.apache.kafka.streams.kstream.ForeachAction
(apply [_ k v]
(println "处理" k v))))
8.3 取舍
用 Jackdaw 手写:
优点:完全可控、易于理解、无额外拓扑
缺点:有状态/窗口/join 要自己实现
用 Kafka Streams:
优点:内置状态、窗口、join、EOS
缺点:Clojure interop 繁琐、调试复杂
心法:「无状态处理用 Jackdaw,有状态聚合用 Streams」——手写无状态消费简单直接;一旦需要窗口、join、状态存储,Kafka Streams 的成熟度远超自己造轮子。
9. 生产运维
9.1 关键配置
;; 生产者:不丢不重
{:acks :all
:enable.idempotence true
:retries 2147483647
:max.in.flight.requests.per.connection 5}
;; 消费者:可预测
{:enable.auto.commit false ;; 手动提交
:max.poll.records 500 ;; 单批大小
:max.poll.interval.ms 300000} ;; 处理超时上限
9.2 监控指标
| 指标 | 含义 | 告警阈值 |
|---|---|---|
| consumer lag | 消费延迟 | 持续增长 |
| rebalance rate | 再均衡频率 | 频繁 |
| produce error rate | 生产失败率 | > 0 |
| request latency | 请求延迟 | 突增 |
9.3 常见坑
| 坑 | 现象 | 规避 |
|---|---|---|
| max.poll.interval 太小 | 反复再均衡 | 调大或减小批次 |
| 自动提交 + 长处理 | 位移超前,丢消息 | 手动提交 |
| key 为 null | 分区不均、无序 | 指定 key |
| 无 DLQ | 毒消息卡分区 | 死信队列 |
| 日志无 trace | 难排查 | 带 request-id |
心法:消费端的稳定性 = 「批次可控 + 位移可控 + 毒消息可控」——批次别太大(避免 poll 超时)、位移手动提交(避免丢)、毒消息进 DLQ(避免卡)。三者做到,消费组就很稳。
10. 速查表与一句话记忆
| 需求 | 写法 |
|---|---|
| 生产者 | kc/producer + kc/produce! |
| 消费者 | kc/consumer + kc/subscribe + kc/poll |
| 管理 topic | ka/create-topics! |
| 分区依据 | key(同 key 同分区) |
| 序列化 | :edn / :json / :avro |
| Schema | Schema Registry |
| 位移提交 | kc/commit-sync |
| 死信队列 | 单独 topic 存毒消息 |
| 背压 | core.async 有界通道 |
| 有状态 | Kafka Streams |
| 不丢不重 | acks=all + 幂等 |
| 监控 | consumer lag |
一句话记忆:Clojure 接 Kafka = Jackdaw(配置即 map、序列化即关键字)→ key 决定分区与顺序 → 序列化用 Schema Registry 立契约(加字段带默认值)→ 位移手动提交、处理完再 commit(at-least-once + 幂等)→ 毒消息进死信队列不卡分区 → core.async 有界通道给消费加背压 → 无状态用 Jackdaw 手写、有状态用 Kafka Streams → 生产靠 acks=all + 幂等 + lag 监控——把事件流接成一条「不丢、不重、不卡、可观测」的主干道。
延伸阅读
- Clojure 微服务架构实战 — 服务间事件解耦
- Clojure 事件溯源与 CQRS — 事件日志与读模型
- Clojure core.async 深入 — 通道与背压编排
- Clojure 数据管道与流处理 — 流式数据处理
- Clojure 并发设计模式 — 并发与消息传递
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。