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

四段拼起来就是完整的故事:

  1. offset 0..5:提交事务的数据批,epoch 0。
  2. offset 6:COMMIT 控制批,epoch 0,序列号 -1。
  3. offset 7..14:被中止事务的数据批,epoch 1(新事务,epoch 递增)。
  4. 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 就完事。

还有两个实际代价:

主动破坏:构造一个「开着不提交」的事务

begin_transaction()、写几条、然后不 commit 也不 abort(进程直接退出)。观察 orders-11 的 LEO 前进但 read_committed 消费者的可读位点停住——因为 LSO 卡在那里。这也解释了为什么事务必须设 transaction.timeout.ms(默认 60000):协调者会在超时后自动中止这个遗留事务,避免它无限期卡住消费者。这个练习能让你看清「精确一次」不是免费的,它把可用性的一部分交给了超时机制。


生产边界

动手

  1. 复现第 7 章的对照:一次 commit、一次 abort,然后用 --isolation-level read_committed 和 read_uncommitted 各读一遍,说出差异条数。
  2. 用 dumplog 在一个有事务输出的分区里找出所有 isControl: true 的批,说出每个控制批标记的是提交还是中止(提示:看它前面那段的 epoch 和它后面的数据可见性)。
  3. 断言题:把一个事务的 send_offsets_to_transaction 去掉只提交输出,重启后为什么会出现重复消费?画出位点与输出的两条时间线。

自测

  1. 幂等生产者和事务都用到 producerId,它们解决的重问题有什么不同?为什么幂等跨不过进程重启?
  2. send_offsets_to_transaction 里为什么要传 consumer_group_metadata,只传 offset 列表行不行?
  3. read_committed 的消费者为什么会因为一个与它无关的长事务而读不到新数据?LSO 在这里是什么角色?
  4. 控制批的 baseSequence 为什么是 -1?如果消费者把控制批当成普通数据返回,会发生什么?
  5. 一套链路是「Kafka → 处理 → 写 MySQL」,你用了事务,能宣称端到端精确一次吗?缺哪一环?

↓ 下一步:06 章 · 容量、积压与运维

进入 keel 阅读