阶段二 · 架构原理与调优

事务与 Exactly-Once:跨分区原子写入

前置回顾

  • 第 03 篇的幂等生产者用 PID 加序号去重,解决的是"单分区、单会话"内的重复
  • 第 06 篇给了事务的第一眼与边界:consume-transform-produce 的 EOS,不覆盖外部系统
  • 本篇拆事务的完整机制:TransactionalID、协调者、两阶段提交、隔离级别,以及 Go 里的工程落地

一句话总结

事务把幂等的保证从"单分区、单会话"升级为跨分区、跨会话,要么全成,要么全败.
机制 = 稳定的 TransactionalID 换 PID + Transaction Coordinator 两阶段提交 + 消费者按隔离级别过滤未提交数据.
边界:只在 Kafka 内部生效,写外部系统要靠幂等处理或 Outbox 模式(阶段三第 20 篇).

幂等解决不了什么

第 03 篇的幂等生产者保证了:一次会话内,同一分区的重试不会产生重复消息.但把视角升到"一次完整的业务处理",还剩三个缺口:

consume-transform-produce 流水线:
1. 从 input-orders 读一条
2. 处理(算价、风控)
3. 写 output-orders
4. 提交 input 的 offset
缺口 场景 后果
跨分区不原子 步骤 3 写两个分区,一个成功一个失败 output 出现"半成品"
跨会话不幂等 步骤 3 成功后、步骤 4 前 crash 重启后重新处理,output 重复
提交不联动 步骤 4 的 offset 提交独立于步骤 3 两个方向的不一致(第 06 篇)

幂等只能管"同一次会话内、同一个分区的重试",三个缺口都在它的边界之外.事务就是为补齐这三个缺口设计的:把"多个分区的写入 + offset 提交"打包成一个原子单元.


事务的三个角色

Producer ──────► Transaction Coordinator(某个 Broker)
│ │
│ │ 事务状态写进 __transaction_state(内部 topic)
│ ▼
├─ 写 output 的 P0, P1 (Broker 侧记录事务标记)
└─ 提交 input offset (第 04 篇的提交, 现在属于事务)
角色 说明
TransactionalID 应用声明的稳定标识,一个实例一个;事务的"账户名"
PID + Epoch 生产者会话标识;带 TransactionalID 时跨会话延续 PID,并把 epoch 加一
Transaction Coordinator 每个 TransactionalID 由某个 Broker 担任协调者,管理事务状态机

内部 topic __transaction_state 又是"内部状态也是日志"的套路(第 08 篇的 __consumer_offsets、第 10 篇的 __cluster_metadata):事务的状态与结果,本身就是一条有序日志.

事务 = 对公转账:个人转账可能出现"扣了款没到账";对公转账要么两边同时记账,要么整笔退回,中间状态谁都查不到.
TransactionalID = 公司账户:操作员换了(实例重启),账户不变,银行认户不认人,还能识别出旧操作员已失效.


事务的完整流程

1. InitTransactions
Producer ──► Coordinator: 注册 TransactionalID, 拿到 PID, epoch 加一
(旧实例若还活着: epoch 落后 → 被 Fencing)

2. BeginTransaction
本地动作, 无需 RPC

3. 写入阶段
首写某个分区前, 先向 Coordinator 登记该分区(AddPartitionsToTxn)
然后 Producer ──► Broker: 写入 output 各分区的数据

4. SendOffsetsToTransaction
把 input 的 offset 提交也纳入本事务

5. 提交 / 回滚(两阶段)
Coordinator 先写 PREPARE_COMMIT 到 __transaction_state
→ 向所有参与分区写 commit marker(控制消息)
→ 写 COMPLETE_COMMIT, 事务结束

两个值得记住的细节:

fencing:同 TransactionalID 的新实例启动时 epoch 加一,旧实例的任何写入都会被 Broker 以 FENCED_INSTANCE_ID 拒绝.滚动发布中新旧实例短暂并存,旧实例自动失效,保证任一时刻只有一个"有效写者".这既是安全机制,也是"实例级唯一 ID"这个约束的由来.

控制消息(commit / abort marker):marker 也是写进分区日志的条目,占用 offset 空间,但对普通消费者不可见(被过滤).回看第 02 篇:offset 只增不减,事务"反悔"用的不是删除,而是"写入标记,让读者跳过".


隔离级别:读者怎么过滤

事务写下的数据,在提交前对消费者可见吗?由消费者的隔离级别决定:

隔离级别 能读到 用途
read_uncommitted(默认) 包括未提交、将被回滚的数据 调试、非事务场景
read_committed 只读已提交的数据 事务型消费者必须用

read_committed 的实现靠 LSO(Last Stable Offset,最后稳定偏移):

Partition 日志:

[普通][事务A: 已提交][事务B: 进行中 ][事务C: 已提交]
▲
LSO 卡在这里

read_committed 消费者最多读到 LSO:
● 事务 A 和它之前的普通消息: 可读
● 事务 B 及其之后(含已提交的 C): 读不到

一个未完成的事务会把 LSO 卡住,后面的已提交数据一起被挡住.所以长事务会阻塞读,这也是 transaction.timeout.ms(默认 60s)存在的意义:超时未提交的事务被协调者强制 abort,LSO 放行.

franz-go 里声明隔离级别:

kgo.FetchIsolationLevel(kgo.ReadCommitted()),

Go 实战:consume-transform-produce

franz-go 提供了 GroupTransactSession,把"事务 + 消费组 + offset 提交进事务"的样板封装好了:

session, err := kgo.NewGroupTransactSession(
kgo.SeedBrokers("localhost:9092"),
kgo.TransactionalID("order-etl-0"), // 每实例唯一, 见下节
kgo.ConsumerGroup("order-etl"),
kgo.ConsumeTopics("input-orders"),

// 读上游也必须是 read_committed, 否则可能读到将被回滚的数据
kgo.FetchIsolationLevel(kgo.ReadCommitted()),
)
if err != nil {
log.Fatal(err)
}

for {
fetches := session.PollFetches(ctx)
if fetches.IsClientClosed() {
return
}
fetches.EachError(func(topic string, partition int32, err error) {
slog.Error("fetch error", "topic", topic, "partition", partition, "err", err)
})
if fetches.NumRecords() == 0 {
continue
}

fetches.EachRecord(func(r *kgo.Record) {
result := transform(r)

// 写入下游: 自动属于当前事务
session.Produce(ctx, &kgo.Record{
Topic: "output-orders",
Key: r.Key,
Value: result,
}, nil)
})

// 提交事务: 下游写入 + 上游 offset 一起原子提交
if err := session.EndTransaction(ctx, kgo.TryCommit); err != nil {
// 提交失败(网络、协调者切换等): 事务被 abort
// 下游无痕, offset 未推进 → 下一轮重新拉取处理, 不产生重复输出
slog.Error("transaction failed", "err", err)
continue
}
}

对照第 06 篇手写 BeginTransaction / EndTransaction 的版本,GroupTransactSession 处理了三件容易写错的事:

  1. offset 提交自动进入事务(不用手工凑 SendOffsetsToTransaction 的时机)
  2. 提交失败自动 abort,事务不留半成品
  3. 下一轮循环自动开启新事务

abort 会重放,处理逻辑要有边界

事务 abort 后,这批消息会被重新拉取、重新处理.transform 里如果有事务外的副作用(发短信、调第三方接口),重放就会重复执行.
原则:副作用要么放进事务覆盖的 Kafka 写入,要么做成幂等操作.


事务的代价与使用约束

项 说明
延迟 每个事务多出两阶段提交与 marker 写入,比非事务慢一截
粒度 按批(每次 Poll 一个事务)摊薄开销;一条一事务会把开销放大到不可用
读侧 read_committed 消费者受 LSO 牵制,长事务卡住整片读
超时 事务必须在 transaction.timeout.ms(默认 60s)内提交完

三条工程约束:

  1. TransactionalID 实例级唯一且稳定:两个实例共用一个 ID 会互相 fencing,循环重启;K8s 上用 Pod 名(与第 11 篇静态成员同一思路)
  2. 事务不覆盖外部系统:数据库、Redis、HTTP 调用不在事务里,这部分靠业务幂等或 Outbox 模式(第 06 篇策略三,阶段三第 20 篇展开)
  3. 只有 consume-transform-produce 模式享受完整 EOS:Kafka 进、Kafka 出的链路可以端到端不丢不重;出到外部系统的链路,事务只能保证 Kafka 侧

生产注意事项

transaction.timeout.ms 要大于单批最长处理时间

单批处理(含下游 RPC)超过超时值,协调者主动 abort 事务,提交必然失败,表现为反复重试但不推进.
批越大这个值越要留足余量,同时盯着端到端延迟.

长事务阻塞读

LSO 被未完成事务卡住,read_committed 的消费者集体等待.
事务粒度保持在"单批"级别,不要跨批次攒大事务.

fencing 是双刃剑

新实例 fencing 旧实例是发布期的安全保证;但配置失误导致两个实例共用一个 TransactionalID 时,它们会互相踢,日志里是反复的 FENCED_INSTANCE_ID.
遇到这个错误,第一反应是查:这个 ID 是不是被两个进程同时用了.

EOS 配置清单

  1. 生产者:幂等默认开启 + TransactionalID(实例唯一)
  2. 消费者:FetchIsolationLevel(ReadCommitted) + 用 GroupTransactSession 承载循环
  3. 事务粒度:每批一个,transaction.timeout.ms 留足处理余量
  4. 边界确认:副作用不出 Kafka;出 Kafka 的部分用幂等或 Outbox 兜底

事务能保证业务逻辑只执行一次吗?

严格说是"输出恰好一次":事务覆盖的 Kafka 写入不重不漏.
处理逻辑本身在 abort 后可能重放,所以逻辑里的外部副作用要幂等,或者搬进事务覆盖的写入里.

事务和幂等生产者是什么关系?

事务建立在幂等之上:事务型生产者先拿 PID(幂等的基础设施),Epoch 机制再加一层跨会话的 fencing.
幂等是"单分区、单会话"的子集保证,事务是它的超集.

两个实例共用一个 TransactionalID 会怎样?

后启动的实例 epoch 自增,前一个实例的写入全部被拒(FENCED_INSTANCE_ID),且会不断交替.
唯一的合法用法是"滚动发布时短暂并存",两个实例不能长期共存.

read_committed 的消费者要等多久?

等 LSO 越过所有未完成事务.正常提交的事务是毫秒级;异常事务最多等到 transaction.timeout.ms 被强制 abort.
这个超时值也是"读侧最坏等待时间"的上界.

写数据库加发 Kafka 怎么做到原子?

事务做不到,协调者管不了外部系统(第 06 篇强调过).
两条路:业务幂等(下游可重复执行),或 Outbox 模式(数据库事务里写 outbox 表,再异步投递 Kafka).后者在阶段三第 20 篇展开.


快速回顾

  • 事务补的缺口:跨分区原子、跨会话幂等、offset 与写入联动,三者都在幂等生产者边界之外
  • 三个角色:TransactionalID(账户)、PID + Epoch(会话)、Transaction Coordinator(协调者),状态存 __transaction_state
  • 两阶段提交:PREPARE_COMMIT → 各分区写 commit marker → COMPLETE;abort 同理写 abort marker
  • fencing:epoch 自增让旧实例失效,是"同 ID 不可并存"与滚动发布安全的来源
  • 隔离级别:read_committed 靠 LSO 过滤未提交数据,长事务会阻塞读
  • 工程落地:GroupTransactSession 封装 consume-transform-produce;事务按批,超时留足余量,副作用不出 Kafka

动手练习

  1. 观察隔离级别:起一个事务生产者,写两条消息后 abort,分别用 read_committed 与 read_uncommitted 消费者各读一遍,对比结果.
  2. console 验证:用 kafka-console-consumer.sh --isolation-level read_committed 观察第 1 步的事务数据被过滤.
  3. 跑通 ETL:用 GroupTransactSession 实现 input → transform → output 最小流水线,在 transform 里故意 panic,验证 abort 后 output 无脏数据、重启后消息被重新处理.
  4. fencing 实验:同一 TransactionalID 启动两个实例,在日志里找到 FENCED_INSTANCE_ID,观察被踢的是哪一个.
  5. 超时实验:把 transaction.timeout.ms 调到 10s,处理逻辑 sleep 15s,观察协调者强制 abort 后提交失败的表现.
  6. 代价对比:同一批数据,事务与非事务各写一遍,对比端到端延迟与吞吐,记录差值.

阶段二(架构原理与调优)七篇至此完成:存储引擎(08)、副本与 ISR(09)、Controller 与元数据(10)、分区再平衡(11)、生产调优(12)、消费调优(13)、事务与 EOS(14)."Kafka 为什么快、为什么可靠、怎么调"的完整图景已经拼齐.

下一阶段(阶段三 · 流处理与高级应用)从第 15 篇开始:跳出"消息队列"的视角,看 Kafka 作为事件流平台的设计哲学,Log 即数据库.