一条 Kafka 记录已经“提交”,不代表业务已经处理;消费者提交了 offset,也不代表外部扣款成功;Kafka 事务提交,更不代表任意 HTTP 服务与数据库都一起提交。三个动作都叫 commit,却属于不同状态机。

分布式系统(28):缓存失效、填充与旧值竞态

第 28 篇把数据库状态、缓存值和填充资格拆开。Kafka 也需要同样的拆分:分区副本决定哪段日志可承诺,consumer group 保存从哪里恢复,transaction coordinator 决定一组 Kafka 写入是否可见。外部副作用若没有加入同一个原子协议,仍在边界之外。

本文用五个可观察量贯穿故障时间线:已复制并提交的日志位置、消费者当前位置、消费组已提交位点、事务状态、外部副作用次数。

四层状态不能合成一个“已处理”

Kafka topic 被分成多个 partition。每个 partition 是有序追加日志,offset 只在该 partition 内有意义。传统 consumer group 把一个 partition 在同一时刻分配给一个成员,成员按 offset 拉取;跨 partition 没有天然总顺序。

flowchart LR
  P[Producer] --> L0[Partition 0 leader]
  P --> L1[Partition 1 leader]
  L0 --> F0[Partition 0 followers]
  L1 --> F1[Partition 1 followers]
  L0 --> C0[Consumer A]
  L1 --> C1[Consumer B]
  C0 --> O[Group committed offsets]
  C1 --> O

需要分别记住四层提交:

层次 决定者 它回答的问题 没有回答的问题
分区日志提交 partition leader 与 ISR 记录是否进入可承诺日志前缀 某个 consumer 是否处理
消费组 offset 提交 group coordinator 成员重启后从哪里继续 处理产生的结果是否成功
Kafka 事务提交 transaction coordinator、相关 partition 与 group coordinator Kafka 输出和组 offset 是否共同可见 外部数据库/HTTP 是否一起提交
外部效果提交 外部系统 扣款、发信、写库等是否发生 Kafka 是否推进 offset

可迁移的判断方式是“先列状态机,再问原子边界”。两个状态由不同服务分别保存,又没有共同提交协议时,任何调用顺序都存在崩溃缝隙。

分区复制只承诺日志前缀

一个 partition 有一个 leader 和若干 follower。producer 写 leader,follower 拉取同一顺序。Kafka 4.3 Design 把记录的 committed 定义为:当前 ISR 中的副本都已经把它追加到日志。acks=all 要求 leader 等待 ISR 确认;它和 min.insync.replicas 一起决定故障时继续写还是拒绝写。[Kafka 4.3 Design][Kafka 4.3 Producer Configs]

sequenceDiagram
  participant P as Producer
  participant L as Leader
  participant F1 as ISR follower 1
  participant F2 as ISR follower 2
  P->>L: produce(offset 8, acks=all)
  L->>F1: replicate offset 8
  L->>F2: replicate offset 8
  F1-->>L: appended
  F2-->>L: appended
  L-->>P: acknowledged
  Note over L,F2: 日志提交位置推进

leader 尾部可以含尚未复制的记录。普通 consumer 不应越过 high watermark 读取这段尾巴,否则换主后可能读到随后消失的数据。acks=1 则允许 leader 本地追加后立即回复;若 leader 在 follower 复制前故障,该条已回复记录可能丢失。即使使用 acks=all,强制从非同步副本做 unclean leader election 也会改变持久性边界。

网络超时仍留下结果未知。producer 没收到回复,无法仅凭超时判断记录是在提交前失败还是提交后丢了响应。Kafka 0.11 起的 idempotent producer 为每个 producer/partition 使用标识与序列号,让受支持的重试不再追加重复记录。它解决的是生产请求重试,不会替 consumer 处理业务,也不会回滚外部系统。

position 与 committed offset 指向不同时间

consumer 的 position 通常表示下一次要返回的 offset,会随 poll 向前。committed group offset 是保存到 Kafka、供分区重分配或进程重启使用的恢复点。日志 high watermark 是 broker 复制状态;三者可以同时不同。

flowchart LR
  H[日志已提交到 offset 9] --> P[本地 position = 8]
  P --> G[组 committed offset = 5]
  G --> R[崩溃恢复从 5 开始]

设输入记录 offset 0 已经复制并提交,group offset 初始为 0。消费者取出记录后,position 变成 1。若业务处理完成、offset 还没提交就崩溃,新成员仍从 0 读取,业务会执行第二次。这是 at-least-once 常见窗口。

sequenceDiagram
  participant K as Kafka input
  participant C as Consumer
  participant X as External system
  K-->>C: record 0,position=1
  C->>X: effect succeeds,count=1
  Note over C: crash before offset commit
  K-->>C: restart from committed offset 0
  C->>X: effect succeeds again,count=2

把顺序倒过来会得到另一种错误。消费者先提交 group offset 1,再执行外部效果;若两者之间崩溃,恢复从 1 开始,offset 0 不再重放,外部效果却从未发生。Kafka 设计文档把这两种顺序分别归为可能重复与可能漏处理的路径。[Kafka 4.3 Design]

offset 自动提交只改变提交时机,不消除边界。它可能在应用真正完成处理之前推进恢复点。关闭自动提交也只是把选择权交给应用,仍需决定业务结果与 offset 如何协调。

这个模式可称为“恢复游标不是处理凭证”:游标只决定重放起点。听到“消息已消费”,应继续追问说的是本地 position、group committed offset,还是业务结果。

high watermark 与 LSO 也不是一回事

Kafka 事务允许记录先进入分区日志,再由 commit 或 abort marker 决定结果。high watermark 可以已经越过这些物理记录;read_committed consumer 还要受 last stable offset(LSO)限制。存在开放事务时,LSO 停在最早未完成事务之前,后面即使已有复制记录也暂不返回。

flowchart LR
  O0[0: normal] --> O1[1: tx A record]
  O1 --> O2[2: normal]
  O2 --> HW[high watermark]
  LSO[last stable offset] -.停在开放事务前.-> O1

事务提交后,read_committed 可返回其中的记录;事务中止后,这些物理记录仍在日志里,但 consumer 根据控制记录过滤它们。默认 read_uncommitted 会返回中止事务和开放事务中的数据。当前官方 consumer 配置还明确指出:read_committed 最多读到 LSO,而不是 high watermark。[Kafka 4.3 Consumer Configs]

因此,“记录复制成功”与“事务记录对 read_committed 可见”是两件事。长时间开放的事务会拖住 LSO,使后续记录无法向这类消费者交付。这是正确性换来的活性成本,并不表示 follower 没有复制数据。

Kafka 事务怎样闭合 consume-transform-produce

幂等 producer 防单个 producer session 的重试重复,事务则把多个 Kafka partition 的写入合成一次决定。应用提供稳定的 transactional.id;transaction coordinator 管理 producer id、epoch 与事务状态,旧 generation 会被 fencing。事务涉及的 partition 收到 commit 或 abort control marker,consumer 据此解释已经存在的记录。[KIP-98]

consume-transform-produce 还差输入进度。consumer group offset 本身存放在 Kafka 内部 topic,因此 producer 可以用 sendOffsetsToTransaction 把输出记录与待提交 offset 放入同一事务。

sequenceDiagram
  participant C as Consumer
  participant P as Transactional producer
  participant T as Transaction coordinator
  participant O as Output partitions
  participant G as Group coordinator
  C->>P: input 0 / next offset 1
  P->>T: begin transaction
  P->>O: output-0 (transactional)
  P->>G: offset 1 (transactional)
  P->>T: commit
  T->>O: commit marker
  T->>G: commit marker
  Note over O,G: output 与恢复位点共同生效

若事务 abort,read_committed 不返回输出,group committed offset 也不前进。恢复后输入记录再次处理。若事务 commit,两者共同生效,恢复不会再从旧输入开始制造第二份 Kafka 输出。这就是 Kafka 范围内 exactly-once processing 的核心闭包:输入和输出都在 Kafka,消费者使用 read_committed,offset 通过同一个事务提交,旧 producer 与旧 consumer generation 能被 fencing。

它不保证一次 poll 会把跨多个 partition 的整个事务作为单个批次交给 consumer。KIP-98 明确列出限制:consumer 可能只订阅事务涉及的一部分 partition,压缩 topic 还可能让事务中的旧 key 版本被后来记录覆盖。原子可见性不能改写成“应用一次收到完整事务”。

重平衡与旧实例需要两道 fencing

仅靠 producer epoch 还不够处理 consumer group 的动态分配。成员 A 暂停太久后可能被移出 group,partition 转给 B;A 恢复后若还能提交旧任务的 offset,就可能覆盖 B 的进度。KIP-447 让事务性 offset commit 携带 group generation、member id 等元数据,使 group coordinator 能拒绝旧成员。

sequenceDiagram
  participant A as Consumer A / old generation
  participant G as Group coordinator
  participant B as Consumer B / new generation
  A->>G: poll partition P
  Note over A: pause beyond session window
  G->>B: assign P, generation g+1
  A->>G: transactional offset commit with g
  G-->>A: fenced / commit rejected

KIP-447 是协议演进记录,不应把其中早期 exactly_once_beta 配置和升级步骤当作 Kafka 4.3 的当前操作指南。它仍揭示了一条通用规则:数据写入者的代际和工作所有权的代际必须一起验证。只 fencing producer,没有验证 consumer assignment,旧任务仍可能以合法 producer 身份提交过期工作。

Kafka 4.0 又启用了 KIP-890 的 Transactions Server-Side Defense。旧协议曾允许迟到 produce 或 EndTxn 请求跨过事务边界,造成 hanging transaction,严重时破坏 EOS。新版协议为每次事务推进 producer epoch,并把 partition 加入事务的校验移到服务端。Kafka 4.3 文档说明 transaction.version=2 自 4.0 起由服务端自动启用,但旧客户端与升级过程仍要按兼容条件判断。[Kafka Transaction Protocol][KIP-890]

外部数据库和 HTTP 不会被 marker 回滚

把 output topic 换成支付 API,事务闭包立即断开:支付服务不知道 Kafka transaction coordinator 的决定,Kafka abort marker 也不能撤销已经成功的扣款。

sequenceDiagram
  participant C as Consumer
  participant K as Kafka transaction
  participant X as External HTTP service
  C->>K: begin + stage output
  C->>X: POST succeeds,effect=1
  Note over C,K: process fails / transaction aborts
  K-->>C: offset remains old
  C->>X: replay POST,effect=2

反过来先提交 Kafka offset、再调用外部服务,会在中间崩溃时丢效果。交换顺序不能消除窗口。要扩展保证,外部系统必须参与:例如把业务结果和输入 offset 存进同一个数据库事务;让 sink 参与可恢复的两阶段提交;使用稳定幂等键与去重表;或者先写 outbox,再由可重试 relay 传播。具体方案的保证取决于外部资源提供的原子操作。

Kafka 官方 Design 也只把通用 exactly-once 明确限定为读取、处理并写回 Kafka topic;其他 destination generally requires cooperation。这个限定与第 25 篇的原子提交边界一致:协调者只能决定参加协议的资源。

分布式系统(25):分布式事务的决定、补偿与消息边界

本地实验:把五个状态同时打印出来

实验位于 examples/distributed-systems/kafka29/check.py,只使用 Python 标准库。它是确定性状态模型,不启动 Kafka。

1
2
3
4
5
mkdir -p examples/distributed-systems/.build/kafka29/tmp
export TMPDIR="$PWD/examples/distributed-systems/.build/kafka29/tmp"
export TMP="$TMPDIR" TEMP="$TMPDIR" PYTHONDONTWRITEBYTECODE=1
python3 -B examples/distributed-systems/kafka29/check.py \
--output examples/distributed-systems/.build/kafka29/observations.json

Python 3.12.3 的正式运行中,日志提交位置始终为 1。处理后未提交 offset 的场景在恢复后把外部计数从 1 增至 2;先提交 offset 的场景恢复于 1,外部计数保持 0。Kafka 事务开放时,输出 high watermark 为 1、LSO 为 0,read_committed 看不到输出;commit 后输出和 group offset 同时可见,abort 后二者都不生效。

外部效果反例中,Kafka 事务 abort 后外部计数仍为 1,恢复重放后变为 2。有限模型没有实现 broker、ISR、coordinator、磁盘或网络,不能证明 Kafka 产品实现,也没有性能含义。

完整输出见结构化观察,运行和来源证据见实验证据,验证边界见验证说明。

安全性、活性与工程选择

日志安全依赖 ISR 与 leader election 配置。若仍有包含已提交前缀的合格副本,正常选主不应丢已提交记录;允许 unclean election 是拿持久性换可用性。生产者的超时必须保留“结果未知”,幂等重试也要满足当前客户端与 broker 的协议条件。

事务安全要求 coordinator 状态、相关 partition marker、producer epoch 和 consumer group generation 协同。活性则要求这些协调者和足够副本可达,开放事务最终能 commit、abort 或超时终止。read_committed 因 LSO 被阻塞,是协议选择的直接后果。

外部效果没有统一答案。可接受重复且有幂等键时,先做效果再提交 offset 通常更安全;不能重复且外部数据库能保存 offset 时,可把结果与 offset 放入同一数据库事务;必须写回 Kafka 时,使用 Kafka transaction 与 read_committed。声称 exactly-once 前,必须把事务参与者逐个列出来。

听到的说法 应检查的状态 可采用的机制
“消息已提交” partition 的 high watermark/ISR acks=all、min.insync.replicas、合格选主
“消费者处理完了” position、group committed offset、业务结果 明确提交顺序、幂等处理或同库事务
“Kafka 恰好一次” 输出 topic、offset、isolation、fencing transactional producer、sendOffsetsToTransaction、read_committed
“端到端恰好一次” 所有外部资源是否参加原子边界 外部事务、2PC、幂等键、去重或 outbox

两个推演练习

开放事务为什么会挡住后面的普通记录。 partition 上 offset 5 属于未完成事务 T,offset 6 是普通记录,high watermark 已到 7。read_committed consumer 能否先返回 offset 6?如果 T 超时中止,恢复后会看到什么?

为保持 partition offset 顺序,consumer 不能越过最早开放事务先交付 6。T 中止并写入 marker 后,LSO 才能推进;offset 5 被过滤,offset 6 可以返回。复制进度没有回退,改变的是事务可见边界。

同一个 transactional.id 能否保护支付 API。 旧实例 A 已成功调用支付 API,尚未提交 Kafka 事务时暂停;新实例 B 用相同 transactional.id 启动并 fence A。是否已经避免重复扣款?

没有。fencing 能阻止 A 继续写或提交 Kafka 事务,却不能让支付 API 撤销既有结果。B 从旧 group offset 重放时仍可能再次扣款。支付侧需要稳定幂等键、去重记录,或参加覆盖 offset 与支付结果的共同事务。

工程结论

Kafka 的日志提交、消费位点和事务分别解决持久前缀、恢复起点与 Kafka 内原子可见性。它们可以组合出可靠的 consume-transform-produce,但组合的最外层边界仍是 Kafka。外部系统没有参与决定时,重复或遗漏窗口不会因为配置项里出现 exactly_once 就消失。

下一篇进入 Spark/RDD。Kafka 把重放起点保存在日志位点里;Spark 则用 lineage 记录分区怎样算出来,并在 shuffle 与持久化处形成不同的恢复边界。

参考资料