前置回顾
- 第 03 篇的幂等生产者用 PID 加序号去重,解决的是"单分区、单会话"内的重复
- 第 06 篇给了事务的第一眼与边界:consume-transform-produce 的 EOS,不覆盖外部系统
- 本篇拆事务的完整机制:TransactionalID、协调者、两阶段提交、隔离级别,以及 Go 里的工程落地
一句话总结
事务把幂等的保证从"单分区、单会话"升级为跨分区、跨会话,要么全成,要么全败.
机制 = 稳定的 TransactionalID 换 PID + Transaction Coordinator 两阶段提交 + 消费者按隔离级别过滤未提交数据.
边界:只在 Kafka 内部生效,写外部系统要靠幂等处理或 Outbox 模式(阶段三第 20 篇).
幂等解决不了什么
第 03 篇的幂等生产者保证了:一次会话内,同一分区的重试不会产生重复消息.但把视角升到"一次完整的业务处理",还剩三个缺口:
|
| 缺口 | 场景 | 后果 |
|---|---|---|
| 跨分区不原子 | 步骤 3 写两个分区,一个成功一个失败 | output 出现"半成品" |
| 跨会话不幂等 | 步骤 3 成功后、步骤 4 前 crash | 重启后重新处理,output 重复 |
| 提交不联动 | 步骤 4 的 offset 提交独立于步骤 3 | 两个方向的不一致(第 06 篇) |
幂等只能管"同一次会话内、同一个分区的重试",三个缺口都在它的边界之外.事务就是为补齐这三个缺口设计的:把"多个分区的写入 + offset 提交"打包成一个原子单元.
事务的三个角色
|
| 角色 | 说明 |
|---|---|
| TransactionalID | 应用声明的稳定标识,一个实例一个;事务的"账户名" |
| PID + Epoch | 生产者会话标识;带 TransactionalID 时跨会话延续 PID,并把 epoch 加一 |
| Transaction Coordinator | 每个 TransactionalID 由某个 Broker 担任协调者,管理事务状态机 |
内部 topic __transaction_state 又是"内部状态也是日志"的套路(第 08 篇的 __consumer_offsets、第 10 篇的 __cluster_metadata):事务的状态与结果,本身就是一条有序日志.
事务 = 对公转账:个人转账可能出现"扣了款没到账";对公转账要么两边同时记账,要么整笔退回,中间状态谁都查不到.
TransactionalID = 公司账户:操作员换了(实例重启),账户不变,银行认户不认人,还能识别出旧操作员已失效.
事务的完整流程
|
两个值得记住的细节:
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,最后稳定偏移):
|
一个未完成的事务会把 LSO 卡住,后面的已提交数据一起被挡住.所以长事务会阻塞读,这也是 transaction.timeout.ms(默认 60s)存在的意义:超时未提交的事务被协调者强制 abort,LSO 放行.
franz-go 里声明隔离级别:
|
Go 实战:consume-transform-produce
franz-go 提供了 GroupTransactSession,把"事务 + 消费组 + offset 提交进事务"的样板封装好了:
|
对照第 06 篇手写 BeginTransaction / EndTransaction 的版本,GroupTransactSession 处理了三件容易写错的事:
- offset 提交自动进入事务(不用手工凑
SendOffsetsToTransaction的时机) - 提交失败自动 abort,事务不留半成品
- 下一轮循环自动开启新事务
abort 会重放,处理逻辑要有边界
事务 abort 后,这批消息会被重新拉取、重新处理.transform 里如果有事务外的副作用(发短信、调第三方接口),重放就会重复执行.
原则:副作用要么放进事务覆盖的 Kafka 写入,要么做成幂等操作.
事务的代价与使用约束
| 项 | 说明 |
|---|---|
| 延迟 | 每个事务多出两阶段提交与 marker 写入,比非事务慢一截 |
| 粒度 | 按批(每次 Poll 一个事务)摊薄开销;一条一事务会把开销放大到不可用 |
| 读侧 | read_committed 消费者受 LSO 牵制,长事务卡住整片读 |
| 超时 | 事务必须在 transaction.timeout.ms(默认 60s)内提交完 |
三条工程约束:
- TransactionalID 实例级唯一且稳定:两个实例共用一个 ID 会互相 fencing,循环重启;K8s 上用 Pod 名(与第 11 篇静态成员同一思路)
- 事务不覆盖外部系统:数据库、Redis、HTTP 调用不在事务里,这部分靠业务幂等或 Outbox 模式(第 06 篇策略三,阶段三第 20 篇展开)
- 只有 consume-transform-produce 模式享受完整 EOS:Kafka 进、Kafka 出的链路可以端到端不丢不重;出到外部系统的链路,事务只能保证 Kafka 侧
生产注意事项
transaction.timeout.ms 要大于单批最长处理时间
单批处理(含下游 RPC)超过超时值,协调者主动 abort 事务,提交必然失败,表现为反复重试但不推进.
批越大这个值越要留足余量,同时盯着端到端延迟.
长事务阻塞读
LSO 被未完成事务卡住,read_committed 的消费者集体等待.
事务粒度保持在"单批"级别,不要跨批次攒大事务.
fencing 是双刃剑
新实例 fencing 旧实例是发布期的安全保证;但配置失误导致两个实例共用一个 TransactionalID 时,它们会互相踢,日志里是反复的 FENCED_INSTANCE_ID.
遇到这个错误,第一反应是查:这个 ID 是不是被两个进程同时用了.
EOS 配置清单
- 生产者:幂等默认开启 +
TransactionalID(实例唯一) - 消费者:
FetchIsolationLevel(ReadCommitted)+ 用GroupTransactSession承载循环 - 事务粒度:每批一个,
transaction.timeout.ms留足处理余量 - 边界确认:副作用不出 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
动手练习
- 观察隔离级别:起一个事务生产者,写两条消息后 abort,分别用 read_committed 与 read_uncommitted 消费者各读一遍,对比结果.
- console 验证:用
kafka-console-consumer.sh --isolation-level read_committed观察第 1 步的事务数据被过滤. - 跑通 ETL:用 GroupTransactSession 实现 input → transform → output 最小流水线,在 transform 里故意 panic,验证 abort 后 output 无脏数据、重启后消息被重新处理.
- fencing 实验:同一 TransactionalID 启动两个实例,在日志里找到 FENCED_INSTANCE_ID,观察被踢的是哪一个.
- 超时实验:把 transaction.timeout.ms 调到 10s,处理逻辑 sleep 15s,观察协调者强制 abort 后提交失败的表现.
- 代价对比:同一批数据,事务与非事务各写一遍,对比端到端延迟与吞吐,记录差值.
阶段二(架构原理与调优)七篇至此完成:存储引擎(08)、副本与 ISR(09)、Controller 与元数据(10)、分区再平衡(11)、生产调优(12)、消费调优(13)、事务与 EOS(14)."Kafka 为什么快、为什么可靠、怎么调"的完整图景已经拼齐.
下一阶段(阶段三 · 流处理与高级应用)从第 15 篇开始:跳出"消息队列"的视角,看 Kafka 作为事件流平台的设计哲学,Log 即数据库.