KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
05 · 端到端可靠性 — keel 龙骨
「先消费、再处理、再写下游、再提交位点」这四步里,只要第 4 步之前进程挂了,消息就会被处理两次。第 03 章的幂等生产者只覆盖了生产端一个会话,跨不过「消费 + 生产」这个组合。这一章讲 Kafka 用事务把两个写变成一件事。
「先消费、再处理、再写下游、再提交位点」这四步里,只要第 4 步之前进程挂了,消息就会被处理两次。第 03 章的幂等生产者只覆盖了生产端一个会话,跨不过「消费 + 生产」这个组合。这一章讲 Kafka 用事务把两个写变成一件事。
现场
一个流处理任务:从 orders-10 读订单,加工之后写到 orders-11,同时要推进消费位点。它面对的真实故障是:写出结果成功了,但提交位点之前进程崩溃。重启之后,这个消费组会从上次提交的位点重读,把同一批订单再加工一次——下游 orders-11 出现了重复。
把「写结果」和「提交位点」放进同一个事务,就能让它们要么都生效、要么都不生效。这一章把这个动作拆到磁盘字节级别。
全链路:read-process-write 的两个写如何变成一个原子动作
flowchart TD
subgraph IN["输入 topic orders-10(3 分区)"]
I["orders-10-0 / -1 / -2 的订单记录"]
end
subgraph APP["消费者进程(group txg1)"]
C["① poll(isolation.level=read_committed)"]
R["② 拿到 8 条及其 (partition, offset)"]
PR["③ 业务处理"]
TX["④ producer.begin_transaction()"]
W1["⑤ produce(orders-11, 处理结果)"]
W2["⑥ send_offsets_to_transaction(<br/>offset+1 列表, consumer_group_metadata)"]
end
subgraph OUT["输出 topic orders-11(3 分区)"]
O["orders-11-0 等分区的数据 + 控制批"]
end
subgraph OFF["内部 topic __consumer_offsets"]
OF["group txg1 的位点记录"]
end
I --> C --> R --> PR --> TX
TX --> W1 --> O
TX --> W2 --> OF
O --> CMT{"⑦ commit_transaction()<br/>还是 abort_transaction()"}
OF --> CMT
CMT -- "commit" --> OK["⑧ orders-11 写入可见 + 位点前进<br/>两者一起生效"]
CMT -. "abort(失败路径)" .-> NG["⑧ orders-11 写入对 read_committed 不可见<br/>位点不动,重启后重读"]
style I fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
style C fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style R fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style PR fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style TX fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
style W1 fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
style W2 fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
style O fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
style OF fill:#fff9db,stroke:#f08c00,color:#1a1a2e
style CMT fill:#fff9db,stroke:#f08c00,color:#1a1a2e
style OK fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style NG fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e
关键在于 ⑤ 和 ⑥ 之间:位点的写不是单独提交的,它是同一个事务里的第二个写。图里 orders-11(业务结果)和 __consumer_offsets(位点)在同一个事务边界内,提交这一个动作让它们同时可见,中止这一个动作让它们同时回退。
精确定义
至少一次(at-least-once):默认语义。消息不会丢,但可能重复——生产者重试、消费者重复消费都会带来重复。
最多一次(at-most-once):先提交位点再处理,挂了就丢。几乎不用。
精确一次(exactly-once):在 Kafka 内部的读-处理-写链路上,每条输入恰好影响一次输出。它由三样东西拼成:幂等生产者(去重)+ 事务(把多个写和位点绑定)+ read_committed 消费(过滤未提交)。
transactional.id:事务生产者的身份。它绑定了 producerId 和 producerEpoch,broker 靠它在生产者重启后识别「这个逻辑生产者」,并中止上一个未完成的事务——这是幂等生产者跨会话做不到的事。
producerEpoch:每次 init_transactions / 新事务会递增,用来隔离同一 transactional.id 的旧会话遗留的批。本机实验里从 epoch 0 变成了 epoch 1。
isolation.level:消费者读隔离级别。read_uncommitted(默认)读所有已写记录(含未提交、含被中止的);read_committed 只读已提交的,并且只能读到 LSO(last stable offset,即最早的未完成事务之前的位点)。
控制批(control batch):写进日志的一种特殊批,标记事务的 COMMIT 或 ABORT。它的 baseSequence/lastSequence 都是 -1,isControl: true。消费者不会收到它(客户端会过滤),但它是磁盘上事务边界的唯一标记。
证据一:提交与中止两个事务
orders-10 里先放 20 条输入。第一个事务读 8 条、写 orders-11、提交(含位点):
[tx] init_transactions done, transactional.id=tx-rw-1, mode=commit
[tx] 读到 8 条, 位点示例: [(2, 0), (2, 1), (2, 2)]
[tx] committed: 写出 8 条到 orders-11 + 提交 8 个位点
再产 10 条输入,第二个事务读 10 条、写 orders-11、中止:
[tx] init_transactions done, transactional.id=tx-rw-1, mode=abort
[tx] 读到 10 条, 位点示例: [(2, 1), (2, 2), (2, 3)]
[tx] aborted: 丢弃对 orders-11 的 10 条写入, 位点不提交
证据二:两种隔离级别看到的差异
同一个 topic orders-11,同一批数据,只换一个参数:
[read_committed] orders-11: 读到 8 条
[read_uncommitted] orders-11: 读到 16 条
read_committed 只看到提交事务写的 8 条;read_uncommitted 多看到 8 条——那正是被中止事务写进去、但对已提交消费者不可见的数据。同一份日志、同一时刻、两个消费者看到不同的数据集,这就是隔离级别在干的事。
证据三:控制批在磁盘上长什么样
dumplog orders-11 的分区 0(这个分区数据最多):
baseOffset: 0 lastOffset: 5 count: 6 baseSequence: 0 lastSequence: 5
producerId: 4001 producerEpoch: 0 isTransactional: true isControl: false
baseOffset: 6 lastOffset: 6 count: 1 baseSequence: -1 lastSequence: -1
producerId: 4001 producerEpoch: 0 isTransactional: true isControl: true
baseOffset: 7 lastOffset: 14 count: 8 baseSequence: 0 lastSequence: 7
producerId: 4001 producerEpoch: 1 isTransactional: true isControl: false
baseOffset: 15 lastOffset: 15 count: 1 baseSequence: -1 lastSequence: -1
producerId: 4001 producerEpoch: 1 isTransactional: true isControl: true
四段拼起来就是完整的故事:
- offset 0..5:提交事务的数据批,epoch 0。
- offset 6:COMMIT 控制批,epoch 0,序列号 -1。
- offset 7..14:被中止事务的数据批,epoch 1(新事务,epoch 递增)。
- offset 15:ABORT 控制批,epoch 1。
对账:这个分区物理上有 16 条记录(6 数据 + 1 控制 + 8 数据 + 1 控制)。read_committed 消费者看得到 6 条(第一段);read_uncommitted 看得到 14 条数据(6+8,控制批被客户端过滤)。再看另外两个分区:
orders-11-1: 只有 1 条记录, isControl: true (ABORT), 数据 0 条
orders-11-2: offset 0..1 提交批(epoch0) + offset 2 COMMIT 控制批
所以全 topic 算下来:read_committed 看到 6+2=8 条(正是第一个事务写出的 8 条),read_uncommitted 看到 6+8+2=16 条(多了被中止的 8 条)。这个对账能完全对上,说明隔离级别的差异就是那一段 ABORT 控制批划出的范围。
顺便注意 orders-11-1:它只有一条 ABORT 控制批,没有数据。这说明事务会向所有参与的分区写控制标记,即使那个分区没有数据——因为「这个事务在这个分区上没有输出」这件事本身也要被标记,否则消费者无法确定边界。
证据四:事务路径多了一跳协调
第一次 init_transactions 时客户端打了一串日志:
Failed to acquire transactional PID from broker TxnCoordinator/2: Broker: Not coordinator: retrying
... (重复约 40 秒)
原因是 transactional.id 的协调者由 __transaction_state 的分区 leader 担任,客户端首次找到了一个非协调者 broker,要重试到路由正确为止。这批重试不影响正确性,但它说明一个事实:事务路径比普通生产多一跳「找协调者」,冷启动或协调者迁移时会体现为延迟。生产上要把这一段算进事务处理器的启动时间预算。
事务保证了什么、没保证什么
保证的:在 Kafka 内部的读-处理-写链路上,输入与输出是原子的;read_committed 消费者不会看到未提交或已中止的数据;位点和输出一起前进或一起回退。配合幂等生产者,可以做到「输入恰好影响输出一次」。
没保证的:Kafka 之外的部分它管不着。如果处理过程还写了一张 MySQL 表、发了一条 webhook,那部分不在事务里,事务提交和外部副作用之间仍然可能不一致。这类「Kafka + 外部系统」的端到端精确一次需要另外的机制(本地消息表、outbox、分布式事务或幂等消费),不是开个 transactional.id 就完事。
还有两个实际代价:
- LSO 阻塞。
read_committed消费者只能读到 LSO。一个开着的长事务会把 LSO 卡住,让所有read_committed消费者读不到这个分区的新数据,即使这些数据和那个长事务无关。本机orders-11-0里 offset 7..14 那段 ABORT 批就曾经处于「未完成」状态。 - 段删除被拖住。未完成事务涉及的段(
deleteHorizonMs非空)不会被 retention 删除,避免删掉还在事务里的数据。长事务会拖住一个分区末尾段的清理。
主动破坏:构造一个「开着不提交」的事务
begin_transaction()、写几条、然后不 commit 也不 abort(进程直接退出)。观察 orders-11 的 LEO 前进但 read_committed 消费者的可读位点停住——因为 LSO 卡在那里。这也解释了为什么事务必须设 transaction.timeout.ms(默认 60000):协调者会在超时后自动中止这个遗留事务,避免它无限期卡住消费者。这个练习能让你看清「精确一次」不是免费的,它把可用性的一部分交给了超时机制。
生产边界
- 教学替身 vs 真实依赖:本机事务只涉及两个 Kafka topic 和
__transaction_state,都是本地三节点。真实 EOS 链路通常还要接外部数据库或对象存储,那部分的一致性要另做;本机也无法演示跨数据中心的事务延迟。 - 要盯的指标:
transaction.state.log的副本状态、未完成事务数、LSO 与 LEO 的差距(差距长期不为 0 说明有长事务卡着)、事务协调者分区的位置、transaction.timeout.ms触发的中止次数。 - 失败策略:事务处理器的
transactional.id必须每个实例唯一且稳定(比如「应用名-分区号」),不能随机生成——随机 id 会让「接管未完成事务」的能力失效。进程重启后要依赖init_transactions去中止上一个会话的残留,而不是假设没有残留。
动手
- 复现第 7 章的对照:一次 commit、一次 abort,然后用
--isolation-level read_committed和read_uncommitted各读一遍,说出差异条数。 - 用
dumplog在一个有事务输出的分区里找出所有isControl: true的批,说出每个控制批标记的是提交还是中止(提示:看它前面那段的 epoch 和它后面的数据可见性)。 - 断言题:把一个事务的
send_offsets_to_transaction去掉只提交输出,重启后为什么会出现重复消费?画出位点与输出的两条时间线。
自测
- 幂等生产者和事务都用到
producerId,它们解决的重问题有什么不同?为什么幂等跨不过进程重启? send_offsets_to_transaction里为什么要传consumer_group_metadata,只传 offset 列表行不行?read_committed的消费者为什么会因为一个与它无关的长事务而读不到新数据?LSO 在这里是什么角色?- 控制批的
baseSequence为什么是 -1?如果消费者把控制批当成普通数据返回,会发生什么? - 一套链路是「Kafka → 处理 → 写 MySQL」,你用了事务,能宣称端到端精确一次吗?缺哪一环?
↓ 下一步:06 章 · 容量、积压与运维