6146 字
31 分钟
Kafka 可靠性与事务:从副本、ISR 到 Exactly Once 的 Failure Model

Kafka 可靠性与事务:从副本、ISR 到 Exactly Once 的 Failure Model#

深入 acks、High Watermark、ELR、幂等 Producer、Transaction Coordinator 与端到端一致性边界

阅读目标#

前四篇已经完成 Kafka 核心数据链的搭建:

第一篇
Log → Partition → Offset → Consumer Position
第二篇
RecordBatch → Segment → Index → Page Cache → Broker Data Path
第三篇
KafkaProducer.send()
→ RecordAccumulator
→ Sender
→ NetworkClient
→ Partition Leader
第四篇
poll()
→ Fetch
→ Consumer Group
→ Rebalance
→ Offset Commit

第五篇要面对一个更现实的问题:

如果链路中的任何节点都可能失败,Kafka 如何定义一条消息究竟是成功、失败、可见、可恢复,还是需要重试?

真实链路中可能发生:

Producer 重复调用 send
Request 没有到达 Broker
Leader 已追加但 ACK 丢失
Follower 落后或掉出 ISR
Leader 在返回成功后永久故障
旧 Leader 在网络分区后恢复
Consumer 读到未完成事务的数据
业务数据库成功,但 Offset Commit 失败

完成本篇后,你应该能够:

  • 从 Failure Boundary 分析消息丢失、重复与乱序;
  • 准确解释 acks=0acks=1acks=all
  • 说明 Replication Factor、Leader、Follower 和 Replica Fetch 的关系;
  • 区分 Replica Set、ISR、ELR、HW、LEO 和 LSO;
  • 解释 min.insync.replicasacks=all 如何共同工作;
  • 解释 Kafka 4.x ELR 与 Unclean Leader Election 的差异;
  • 说明 PID、Producer Epoch 与 Sequence Number 如何实现 Broker 端去重;
  • 区分 Idempotent Producer 与 Kafka Transaction;
  • 画出 Transaction Coordinator 与 __transaction_state
  • 解释 Commit/Abort Marker、LSO 与 read_committed
  • 实现 Kafka → Process → Kafka 的事务处理;
  • 明确 Exactly Once 在 Kafka、数据库和 HTTP 边界上的差异;
  • 根据副本、Producer 和 Transaction 指标排查故障;
  • 沿 ReplicaManager、Partition、ProducerStateManager 与 Transaction Coordinator 阅读源码。

本文以 Apache Kafka 4.3.x 为基线。Kafka 4.3 Producer 默认 acks=all,在不存在冲突配置时默认启用 Idempotence;启用 Idempotence 要求 acks=allretries>0max.in.flight.requests.per.connection<=5。Kafka 4.x 还引入并继续完善 Eligible Leader Replicas(ELR),用于在严格 min.insync.replicas 语义下改进安全选举。


0. Broker 已经写入,但 ACK 丢了:成功了吗#

Producer
↓ ProduceRequest(A)
Partition Leader
↓ A 已追加到 Log
ProduceResponse
↓ ACK 在网络中丢失
Producer Timeout

ACK 丢失后的不确定结果

Producer 无法区分:

Request 根本没有到达 Broker
Broker 收到但尚未追加
Broker 已追加,ACK 在回程丢失
Broker 已追加并复制,但连接中断

如果不重试,第一种情况会丢消息;如果重试,第三种情况可能重复。

这说明:

分布式发送最难的不是服务器明确返回失败,而是客户端无法判断服务器执行到了哪一步。

Idempotent Producer 没有让网络变可靠。它解决的是:

当 Producer 因不确定结果重试同一个 Batch 时,Broker 能否识别它已经追加过。


1. 一条消息的 Failure Boundary#

消息链路中的 Failure Boundary

“Kafka 会不会丢消息”不是一个足够精确的问题。至少要拆成七段。

1.1 Application Boundary#

应用可能因为 HTTP 重试、任务重放或代码错误重复调用:

producer.send(event);

Kafka Idempotence 无法理解两个不同 send() 在业务上是不是同一事件。业务层仍需要:

Event ID
Request ID
订单号 + 事件类型
业务状态机
幂等记录

1.2 Producer Boundary#

Record 可能仍处于:

序列化阶段
RecordAccumulator
等待 Buffer
In-flight Request
Retry Backoff
Delivery Timeout

send() 返回 Future 不等于消息已经进入 Broker。

1.3 Network Boundary#

可能发生:

Request 丢失
Response 丢失
连接中断
Broker 成功但客户端未知

1.4 Leader Boundary#

Leader 负责:

校验 RecordBatch
分配最终 Offset
追加本地 Log
校验 Producer Sequence
等待复制条件
返回 ProduceResponse

必须区分:

Leader 已追加
形成安全复制条件
Producer 已收到成功响应

1.5 Replica Boundary#

Follower 可能:

在线且同步
在线但落后
网络不可达
磁盘过慢
掉出 ISR

三副本不等于当前消息一定已经在三台机器上安全存在。

1.6 Consumer Boundary#

Consumer 是否可见受到:

High Watermark
Isolation Level
Last Stable Offset
Position

限制。

1.7 External Side Effect Boundary#

消费后可能:

更新 MySQL
调用支付 API
修改库存
执行 Agent Tool

这些副作用不属于 Kafka Log,Kafka Transaction 不能自动覆盖它们。

因此排障时固定问:

失败发生在哪一段?
这一段的 Source of Truth 是什么?
调用方能否确定执行结果?
Retry 是否会产生重复?
恢复状态由谁保存?

2. acks:Producer 何时认为写入完成#

acks 配置的确认边界

acks 定义 Producer 等待的 Broker 确认条件。它不是可靠性的唯一开关,可靠性还取决于:

Replication Factor
ISR
min.insync.replicas
Leader Election
Retry
Idempotence
故障域分布
应用是否检查 Future

2.1 acks=0#

acks=0

Producer 不等待 Broker Response。消息进入 Socket Buffer 后就认为已发送。

结果:

Broker 是否收到:未知
是否被拒绝:未知
Retry 无法依赖 Broker 错误
RecordMetadata.offset=-1

只适合可容忍丢失的低价值遥测。

2.2 acks=1#

acks=1

Leader 本地追加后立即返回,不等待 Follower。

Leader 追加 A
→ 返回成功
→ Follower 尚未复制
→ Leader 永久故障
→ 新 Leader 没有 A

Producer 已经看到成功,但 Record 可能丢失。

2.3 acks=all#

acks=all

等价于 -1。Leader 等待当前 ISR 全部确认,并检查 min.insync.replicas

Kafka 4.3 Producer 默认使用 acks=all

2.4 acks=all 不是固定等待两个副本#

假设:

RF=3
ISR={1,2,3}
min ISR=2

acks=all 等待 1、2、3,而不是只等两个。

如果:

ISR={1,2}

等待 1、2。

如果:

ISR={1}

因为不满足 min ISR=2,写入被拒绝。

2.5 成功响应不等于每块磁盘都同步 fsync#

需要区分:

Producer Response Success
Consumer Visibility
Physical Device Flush

Kafka 的核心可靠性依赖副本协议、确认条件和选举约束,不是每条消息都同步刷所有物理磁盘。


3. Kafka 副本模型:Leader 与 Follower#

一个 Partition 的:

Replication Factor=N

表示它配置了 N 个 Replica,其中一个 Leader,其余 Follower。

Follower 主动 Fetch Leader

3.1 Producer 只写 Leader#

所有写入经 Leader 排序并分配 Offset,保证 Partition 只有一个合法追加顺序。

3.2 Follower 主动 Fetch#

Follower LEO
↓ FetchRequest(offset=LEO)
Leader
↓ 返回 RecordBatch
Follower 追加本地 Log
↓ 再次 Fetch

使用 Pull 的好处:

  • 每个 Follower 根据自身进度拉取;
  • 支持批量复制;
  • 慢 Follower 不直接占用 Producer 线程;
  • Leader 可以依据 Fetch 进度判断同步状态。

3.3 Follower 不是独立可写副本#

Kafka 不是多主写。Follower 必须收敛到当前 Leader Epoch 对应的合法历史。

Leader 切换后,旧副本可能需要:

比较 Leader Epoch
定位 Diverging Offset
截断不合法尾部
重新 Fetch

副本系统的目标不是保留每台 Broker 曾经写过的所有尾部,而是形成唯一合法 Partition History。


4. ISR:当前安全同步集合#

ISR 是:

In-Sync Replicas

包括 Leader 和当前满足同步条件的 Followers。

ISR、HW 与 LEO

4.1 ISR 是动态集合#

Replicas=[1,2,3]
ISR=[1,2,3]

如果 Broker 3 长时间不 Fetch 或无法追上:

Replicas=[1,2,3]
ISR=[1,2]

恢复并追上后可重新进入 ISR。

4.2 同步不是零延迟#

Follower 不必每个瞬间与 Leader Offset 完全相同。Kafka 关注它是否在限定时间内持续 Fetch 和追赶。

Kafka 4.3 默认:

replica.lag.time.max.ms=30000

长时间不 Fetch 或不能追上会被移出 ISR。

4.3 ISR 参与哪些机制#

acks=all 的确认范围
min ISR 写入门槛
High Watermark 推进
Leader Election
故障告警

4.4 Replica Set、ISR 与 ELR#

Replica Set
→ 配置的所有副本
ISR
→ 当前同步集合
ELR
→ 已不在 ISR,但仍包含全部安全提交历史、
可以安全参与选举的 Eligible Replicas

5. High Watermark、LEO 与可见性#

5.1 Log End Offset#

每个 Replica 都有自己的 LEO:

Leader LEO=12
Follower A LEO=10
Follower B LEO=8

表示本地下一写入位置。

5.2 High Watermark#

HW 是 Consumer 可安全读取的复制提交边界。

Leader 本地 Offset 8~11 即使存在,也可能还不能返回给 Consumer。否则 Consumer 读到了 11,而故障后新 Leader 只有 8,就出现已经被观察过的数据从合法历史消失。

5.3 HW 不能只背成 min(ISR LEO)#

可用近似模型理解 HW 与最慢同步副本有关,但实现还涉及:

Follower Fetch 确认轮次
Leader Epoch
ISR 变化
严格 min ISR 语义
ELR

不要把简化公式当成完整实现。

5.4 为什么需要下一轮 Fetch#

Follower 请求 offset=5,Leader 返回到 10。Leader不会立即收到独立的磁盘 ACK。

Follower 下一次发:

Fetch(offset=10)

Leader 才知道它已经复制到 10。因此 HW 推进天然包含 Replica Fetch 协议轮次。

5.5 事务下的读取边界#

isolation.level=read_uncommitted

通常读到 HW。

isolation.level=read_committed

还受 LSO 限制,并跳过 Abort Transaction 数据。


6. min.insync.replicas:持久性与可用性的选择#

min ISR 的可用性与持久性权衡

假设:

RF=3
min ISR=2

6.1 ISR={1,2,3}#

满足写入条件,acks=all 等待三个当前 ISR。

6.2 ISR={1,2}#

仍满足 min ISR=2,继续写入,等待两个 ISR。

6.3 ISR={1}#

写入被拒绝,Producer 可能收到:

NotEnoughReplicasException
NotEnoughReplicasAfterAppendException

6.4 为什么拒绝写入#

只剩 Leader 仍继续确认:

唯一 Leader 写入
→ Producer 成功
→ Leader 永久丢失
→ 无副本可恢复

拒绝写入是在明确表达:

当前集群无法继续满足声明的持久性级别。

6.5 权衡#

提高 min ISR:

+ 已确认数据更难丢失
- 故障时更早停止写入

降低:

+ 故障期间更容易保持可写
- 可安全恢复副本更少

常见关键业务组合:

RF=3
min ISR=2
acks=all
unclean election disabled
Rack Awareness

7. Leader Election:ISR、ELR 与 Unclean#

Kafka 4.x Leader Election 候选层级

7.1 优先 ISR#

ISR 中的 Replica 包含安全提交历史,因此首先从 ISR 选择 Leader。

7.2 ELR 为什么出现#

严格 min ISR 下,当 ISR 数低于门槛:

HW 停止推进
新写入无法形成成功确认

此时某些 Replica 即使被移出 ISR,仍可能包含全部已经提交的数据。Kafka 4.0 引入 ELR,Kafka 4.1 起新集群默认启用,用于记录这些仍可安全选举的 Replica。

Kafka 4.3 的候选顺序可以概括为:

ISR 非空 → 从 ISR 选择
否则 ELR 非空 → 从未 Fenced ELR 选择
否则考虑未 Fenced 的 Last Known Leader

7.3 ELR 与 Unclean 的根本差异#

ELR
→ 虽不在 ISR,但仍被证明包含安全提交历史
Unclean
→ 可能缺失已确认数据

ELR 目标是在安全前提下改善可用性;Unclean 用数据风险换可用性。

7.4 Unclean Leader Election#

unclean.leader.election.enable=true

假设:

旧 Leader: 0..100
Follower: 0..80

Follower 成为新 Leader 后,81~100 会从合法历史消失。

关键交易通常不应开启 Unclean Election。

7.5 Leader Epoch 与旧 Leader Fencing#

网络分区后,旧 Leader 可能仍认为自己是 Leader。KRaft Controller 通过 Leader Epoch 和 Metadata 隔离旧 Leader。

旧 Broker 恢复后必须:

确认新 Epoch
截断不合法尾部
作为 Follower 重新追赶

8. Retry、重复与乱序#

Retry 引发的重复与乱序

8.1 ACK 丢失导致重复#

A 已追加
ACK 丢失
Producer Retry A

没有去重状态时,Broker 无法判断它是新 Batch 还是旧 Batch 重试。

8.2 多 In-flight 导致乱序#

同一 Partition:

A 失败等待重试
B 后发先成功
A 重试成功

可能写成:

B → A

Kafka 4.3 官方配置明确说明:禁用 Idempotence、允许 Retry 且 max.in.flight>1 时可能重排。

8.3 把 In-flight 设为 1#

能避免上述重排,但降低连接 Pipeline 吞吐。

Idempotence 允许 max.in.flight<=5 仍保持顺序,因为 Broker 能校验 Sequence Number。

8.4 应用级重发仍可能重复#

即使 Idempotence 开启:

catch (Exception e) {
producer.send(rebuildEvent());
}

第二次调用可能成为新的业务发送,不再是原 Batch Retry。最终失败后是否重发,需要业务 Event ID 和状态查询判断。


9. Idempotent Producer:PID、Epoch 与 Sequence#

Idempotent Producer 的去重状态

Kafka 4.3 在没有冲突配置时默认:

enable.idempotence=true

依赖:

acks=all
retries>0
max.in.flight.requests.per.connection<=5

9.1 Producer ID#

Broker 为 Producer 分配 PID,RecordBatch Header 携带 PID。

9.2 Producer Epoch#

PID 对应 Epoch。新会话取代旧会话时 Epoch 增加,旧 Epoch 会被识别为 Stale。

Transactional Producer 使用同一 transactional.id 重启时,新实例可以 Fence 旧实例。

9.3 Sequence Number#

Producer 为每个:

PID + TopicPartition

维护独立 Sequence。

P0: 0,1,2...
P1: 0,1...

9.4 Broker 的判断#

期望下一批 Seq=3

收到 3 → 正常追加
收到 2 → 可能是已成功 Batch 的 Retry
收到 5 → 中间存在 Gap,状态异常
旧 Epoch → Stale / Fenced

9.5 去重作用域#

Idempotence 能解决:

同一 Producer Session
同一 TopicPartition
同一 Batch 的 Broker Retry

不能解决:

两个 Producer 发送相同 Event
应用重启后重新构造事件
Consumer 重复更新数据库
人工 Replay
HTTP 上游重复请求

9.6 为什么要求 acks=all#

如果使用 acks=1,Leader 确认 Sequence 后永久故障,而 Follower 未复制 Batch 和 Producer State,新 Leader 无法可靠延续去重状态。

9.7 Broker 状态恢复#

源码重点:

ProducerStateManager
ProducerAppendInfo
BatchMetadata
VerificationStateEntry

负责 Sequence 校验、Epoch、Duplicate Detection、Transaction 状态和 Snapshot 恢复。


10. Idempotence 为什么不等于 Transaction#

幂等和事务的边界

Idempotence 解决单 Partition Retry 去重。

业务可能需要:

P0: OrderCreated
P1: InventoryReserved

如果 P0 成功、P1 失败,即使两个写入各自幂等,业务仍是部分成功。

Kafka Transaction 提供:

多个 Topic
多个 Partition
Consumer Offset

的原子可见性。

10.1 事务不是磁盘回滚#

事务 Record 会先追加进普通 Log。Abort 时不物理删除,而是使用:

Transaction Metadata
Abort Marker
Consumer Isolation

read_committed Consumer 跳过。

Kafka 事务控制的是原子可见性,不是 Undo 写入。


11. Transactional Producer API#

配置:

transactional.id=fulfillment-worker-01

设置后自动启用 Idempotence。

11.1 initTransactions()#

producer.initTransactions();

它会:

  1. 查找 Transaction Coordinator;
  2. 完成或 Abort 同 ID 的旧事务;
  3. 获取 PID 与 Epoch;
  4. Fence 旧 Producer 实例。

11.2 beginTransaction()#

producer.beginTransaction();

一个 Producer 同时只能有一个 Open Transaction。

11.3 send()#

beginTransaction() 和 Commit/Abort 之间的所有发送都属于当前事务。

11.4 sendOffsetsToTransaction()#

producer.sendOffsetsToTransaction(
records.nextOffsets(),
consumer.groupMetadata()
);

这些 Offset 只有在事务最终 Commit 后才算已提交。

11.5 commitTransaction()#

会先 Flush 事务内 Record,再提交。事务中任一不可恢复 Send 失败都会导致 Commit 失败。

11.6 abortTransaction()#

将当前事务标记为 Abort,read_committed 不返回其中的业务 Record。

11.7 Transactional ID 设计#

应该稳定且唯一映射一个有状态实例或 Shard:

application + task/shard identity

多个活跃实例共享 ID 会互相 Fence;每次完全随机又无法安全接管旧事务。


12. Transaction Coordinator 与状态机#

Transaction Coordinator 状态机

12.1 Coordinator 定位#

hash(transactional.id)
→ __transaction_state Partition
→ 该 Partition Leader
→ Transaction Coordinator

Kafka 4.3 默认:

transaction.state.log.num.partitions=50
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=2

12.2 Transaction Metadata#

概念上包括:

transactional.id
PID
Producer Epoch
Transaction Timeout
State
参与 TopicPartitions
Start Timestamp

12.3 状态机#

Empty
→ Ongoing
├→ PrepareCommit → CompleteCommit
└→ PrepareAbort → CompleteAbort

完成过程:

  1. 持久化即将 Commit/Abort 的决定;
  2. 向参与 Partition 写 Control Marker;
  3. 所有 Marker 完成后写最终状态。

12.4 AddPartitionsToTxn#

Producer 首次向某 Partition 写事务数据时,需要将其加入当前事务参与集合,因为 Coordinator 最终必须知道向哪些 Partition 写 Marker。

12.5 Transaction Timeout#

Broker 默认上限:

transaction.max.timeout.ms=900000

长事务会:

阻挡 LSO
延迟 read_committed
增加恢复成本
占用 Transaction State

Kafka Transaction 不适合包裹分钟级人工审批或长时间 Agent Workflow。


13. Commit/Abort Marker、LSO 与隔离读取#

事务标记、LSO 与 Consumer 隔离

13.1 事务数据先追加#

事务 Record 已占用正常 Offset。

13.2 Control Marker#

Commit 时写:

COMMIT Marker

Abort 时写:

ABORT Marker

Control Batch 在 Log 中存在,但普通 Consumer 不把它作为业务 Record 返回。

13.3 Last Stable Offset#

LSO 可以理解为最早未完成事务之前的稳定读取边界。

即使 HW 已到 20,最早 Open Transaction 从 4 开始,read_committed 仍可能停在 4 之前。

13.4 Abort Transaction Index#

Broker 使用 Transaction Index 等结构记录中止事务范围。read_committed Fetch 会跳过这些 Record。

13.5 Open Transaction 造成的 Lag#

如果:

HW/LEO 持续前进
read_committed Position 停滞

可能不是 Consumer 慢,而是 LSO 被长事务阻塞。


14. Kafka → Process → Kafka 的 Exactly Once#

Kafka Read-Process-Write 的事务边界

不用事务时:

输出成功
Offset Commit 失败
→ 重启后输出重复

或者:

Offset Commit 成功
输出失败
→ 输入被跳过

14.1 输出和 Offset 进入同一事务#

producer.beginTransaction();
for (ConsumerRecord<String, String> record : records) {
producer.send(transform(record));
}
producer.sendOffsetsToTransaction(
records.nextOffsets(),
consumer.groupMetadata()
);
producer.commitTransaction();

失败:

producer.abortTransaction();

输出不可见,输入 Offset 不推进。

14.2 下游 Consumer#

isolation.level=read_committed

否则可能看到未 Commit 数据。

14.3 Rebalance 与 Fencing#

如果处理期间失去 Partition,sendOffsetsToTransaction() 可能因 Group Metadata 失效而失败。应用应 Abort,重新加入组并从 Committed Offset 恢复。


15. Exactly Once 的边界#

Exactly Once 在不同系统边界上的机制

“Kafka 支持 Exactly Once”必须继续问:

在哪个边界?
对哪个副作用?
由谁观察?

15.1 Producer → Kafka#

Idempotent Producer 保证 Broker Retry 不会在单 Partition Log 中追加重复 Batch。

15.2 Kafka 多 Partition#

Kafka Transaction 提供多个 TopicPartition 原子可见性。

15.3 Kafka Input → Kafka Output#

Transaction
+ sendOffsetsToTransaction
+ read_committed

实现 Kafka 内部 Read-Process-Write EOS。

15.4 Kafka → MySQL#

Kafka Transaction 无法原子提交 MySQL Local Transaction。

MySQL Commit 成功
进程崩溃
Kafka Offset 未 Commit
→ 重复消费

需要:

Event ID 唯一约束
条件更新
Inbox Pattern
业务幂等

业务 DB 更新后发布 Kafka Event,通常使用:

Transactional Outbox
+ CDC / Reliable Publisher

15.5 Kafka → HTTP / Payment#

需要下游支持:

Idempotency-Key
Request ID
状态查询
业务状态机
补偿与对账

15.6 Effectively Once#

跨系统无法建立全局事务时,常见目标:

At-Least-Once Delivery
+ Idempotent Side Effect
= Effectively Once

16. 丢失、重复与乱序的来源#

16.1 Producer 端“丢失”#

  • acks=0
  • 不检查 Future/Callback;
  • 进程退出前未 Close;
  • Delivery Timeout;
  • Buffer 异常被吞;
  • Serializer 失败;
  • 把发送失败误当成功。

16.2 Broker 端“丢失”#

  • RF=1;
  • acks=1 后 Leader 永久故障;
  • min ISR 过低;
  • Unclean Election;
  • 副本位于同一故障域;
  • 错误 Reassignment 或人工操作。

16.3 Consumer 端“丢失”#

消息仍在 Kafka,但业务没执行:

  • Auto Commit 过早;
  • 处理前 Commit;
  • Worker 未完成就推进水位;
  • auto.offset.reset=latest
  • 人工 Reset;
  • Retention 早于恢复窗口。

16.4 重复#

  • ACK 丢失后的非幂等 Retry;
  • 应用重复 Send;
  • 业务成功后 Commit 失败;
  • Rebalance;
  • Replay;
  • Outbox Publisher 重试;
  • 外部 API 无幂等。

16.5 乱序#

Kafka 只保证 Partition 内 Log 顺序。乱序可能来自:

相同 Key 进入不同 Partition
Topic 扩容改变 Key 映射
禁用 Idempotence + 多 In-flight Retry
同 Partition 并发业务 Worker
Retry Topic / DLT
多 Topic 合流
事件时间与到达时间不同

17. 可靠性配置矩阵#

17.1 Producer#

acks=all
enable.idempotence=true
retries=2147483647
max.in.flight.requests.per.connection=5
delivery.timeout.ms=120000
request.timeout.ms=30000

多数场景保留 Kafka 4.3 默认值,主要通过 SLO 调整 delivery.timeout.ms,不要盲目设置巨大 Retry 次数和应用级重发。

17.2 Topic#

关键事件常见:

RF=3
min.insync.replicas=2
unclean.leader.election.enable=false
Rack Awareness

17.3 Transaction Internal Topic#

生产默认:

transaction.state.log.replication.factor=3
transaction.state.log.min.isr=2

单节点实验临时降低后不要复制到生产。

17.4 Consumer#

事务下游:

isolation.level=read_committed
enable.auto.commit=false

具体 Offset 是否手动提交取决于 Kafka Transaction、Streams 或框架 Container Transaction。


18. Metrics 与告警#

Kafka 可靠性故障诊断树

18.1 Broker 副本#

UnderReplicatedPartitions
UnderMinIsrPartitionCount
OfflineReplicaCount
IsrShrinksPerSec
IsrExpandsPerSec
FailedIsrUpdatesPerSec
ReplicaFetcher MaxLag

正常情况下 Offline Replica 为 0,非故障恢复期 ISR Shrink/Expansion 应接近 0。

18.2 Producer#

record-error-rate
record-retry-rate
request-latency
requests-in-flight
record-queue-time
bufferpool-wait-time
produce-throttle-time

18.3 Transaction#

transaction-start-rate
transaction-commit-rate
transaction-abort-rate
transaction-duration
Open Transaction
LSO Lag
Transaction Coordinator Load

实际指标名应以具体 Kafka 4.3.x 客户端和 Reporter 输出为准。

18.4 告警#

至少包括:

UnderMinISR > 0
OfflineReplicaCount > 0
ISR Shrink 突增
Unclean Election
Producer Error/Retry 激增
Transaction Abort 激增
Transaction Duration 超限
read_committed LSO Lag

19. 源码阅读地图#

固定 Kafka 4.3.x Tag 阅读,避免 trunk 方法变化。

19.1 Produce 与 Replica#

KafkaApis
ReplicaManager
Partition
UnifiedLog

关注:

  • ProduceRequest 如何进入 Append;
  • acks 如何影响 Delayed Produce;
  • min ISR 在 Append 前后如何校验;
  • ProduceResponse 何时完成;
  • Leader Epoch 如何传递。

19.2 Replica Fetch#

ReplicaFetcherManager
ReplicaFetcherThread
AbstractFetcherThread

关注:

  • Follower Fetch Offset;
  • FetchResponse 追加;
  • Diverging Epoch;
  • Log Truncation;
  • ISR 恢复。

19.3 ISR、HW 与 ELR#

Partition
ReplicaManager
AlterPartitionManager
KRaft Partition Registration

关注 ISR Expand/Shrink、HW 更新、Controller 确认和 ELR 候选状态。

19.4 Idempotent Producer#

Client:

TransactionManager
ProducerIdAndEpoch
Sender
ProducerBatch

Broker:

ProducerStateManager
ProducerAppendInfo
BatchMetadata
VerificationStateEntry

19.5 Transaction Coordinator#

TransactionCoordinatorService
TransactionCoordinatorShard
TransactionMetadata
TransactionState
TransactionLog
TransactionMarkerChannelManager

关注:

transactional.id 分区
InitProducerId
AddPartitions
Prepare Commit/Abort
Marker
Timeout Abort
Coordinator Failover

19.6 Isolation Fetch#

UnifiedLog.read
FetchIsolation
AbortedTxn
TransactionIndex
CompletedTxn

推荐顺序:

acks/min ISR 文档
→ Partition Append
→ Replica Fetcher
→ ISR/HW/ELR
→ ProducerStateManager
→ Client TransactionManager
→ Transaction Coordinator
→ Marker/LSO/Isolation Fetch

20. 可运行实验:三 Broker KRaft#

以下用于本地学习。故障注入和 Broker Stop 不要直接在生产环境执行。

20.1 Docker Compose#

x-kafka-common: &kafka-common
image: apache/kafka:4.3.0
environment: &kafka-env
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-1:9093,2@kafka-2:9093,3@kafka-3:9093
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT
KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 2
KAFKA_DEFAULT_REPLICATION_FACTOR: 3
KAFKA_MIN_INSYNC_REPLICAS: 2
services:
kafka-1:
<<: *kafka-common
container_name: kafka-r1
hostname: kafka-1
ports:
- "19092:19092"
environment:
<<: *kafka-env
KAFKA_NODE_ID: 1
KAFKA_LISTENERS: INTERNAL://:9092,EXTERNAL://:19092,CONTROLLER://:9093
KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-1:9092,EXTERNAL://localhost:19092
volumes:
- kafka_r1:/var/lib/kafka/data
kafka-2:
<<: *kafka-common
container_name: kafka-r2
hostname: kafka-2
ports:
- "29092:29092"
environment:
<<: *kafka-env
KAFKA_NODE_ID: 2
KAFKA_LISTENERS: INTERNAL://:9092,EXTERNAL://:29092,CONTROLLER://:9093
KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-2:9092,EXTERNAL://localhost:29092
volumes:
- kafka_r2:/var/lib/kafka/data
kafka-3:
<<: *kafka-common
container_name: kafka-r3
hostname: kafka-3
ports:
- "39092:39092"
environment:
<<: *kafka-env
KAFKA_NODE_ID: 3
KAFKA_LISTENERS: INTERNAL://:9092,EXTERNAL://:39092,CONTROLLER://:9093
KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-3:9092,EXTERNAL://localhost:39092
volumes:
- kafka_r3:/var/lib/kafka/data
volumes:
kafka_r1:
kafka_r2:
kafka_r3:

启动:

Terminal window
docker compose up -d

20.2 创建 Topic#

Terminal window
docker exec kafka-r1 \
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server kafka-1:9092 \
--create \
--topic reliable-events \
--partitions 3 \
--replication-factor 3 \
--config min.insync.replicas=2

查看:

Terminal window
docker exec kafka-r1 \
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server kafka-1:9092 \
--describe \
--topic reliable-events

20.3 停止一个 Follower#

使用 acks=all 持续生产,然后停止一个非 Leader Broker:

Terminal window
docker stop kafka-r3

观察 ISR 缩小。RF=3、ISR=2、min ISR=2 时仍可写。

20.4 只剩一个 ISR#

再停止另一个 ISR Broker。如果 Leader 仍在线但 ISR 只剩一个,acks=all 应被拒绝。

恢复:

Terminal window
docker start kafka-r2 kafka-r3

等待 ISR 扩回。

20.5 Leader 切换#

Describe 找到某 Partition Leader,停止该 Broker,观察新 Leader、ISR 和已确认数据是否仍可消费。

20.6 Idempotent Producer#

Properties props = new Properties();
props.put(
ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
"localhost:19092,localhost:29092,localhost:39092"
);
props.put(
ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
StringSerializer.class.getName()
);
props.put(
ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
StringSerializer.class.getName()
);
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");
props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, "5");
try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
for (int i = 0; i < 1000; i++) {
producer.send(
new ProducerRecord<>(
"reliable-events",
"order-" + (i % 10),
"event-" + i
),
(metadata, error) -> {
if (error != null) {
error.printStackTrace();
}
}
);
}
}

20.7 Transactional Producer#

Properties props = new Properties();
props.put(
ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,
"localhost:19092,localhost:29092,localhost:39092"
);
props.put(
ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,
StringSerializer.class.getName()
);
props.put(
ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,
StringSerializer.class.getName()
);
props.put(
ProducerConfig.TRANSACTIONAL_ID_CONFIG,
"transaction-lab-01"
);
try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
producer.initTransactions();
try {
producer.beginTransaction();
producer.send(new ProducerRecord<>(
"reliable-events", 0,
"order-1001", "OrderCreated"
));
producer.send(new ProducerRecord<>(
"reliable-events", 1,
"order-1001", "InventoryReserved"
));
producer.commitTransaction();
} catch (Exception e) {
producer.abortTransaction();
throw e;
}
}

20.8 对比事务可见性#

分别启动:

Terminal window
--consumer-property isolation.level=read_uncommitted

和:

Terminal window
--consumer-property isolation.level=read_committed

将示例改为 abortTransaction(),比较结果。

20.9 Kafka → Kafka 事务模板#

consumer.subscribe(List.of("input-events"));
producer.initTransactions();
while (true) {
ConsumerRecords<String, String> records =
consumer.poll(Duration.ofSeconds(1));
if (records.isEmpty()) {
continue;
}
try {
producer.beginTransaction();
for (ConsumerRecord<String, String> record : records) {
producer.send(new ProducerRecord<>(
"output-events",
record.key(),
transform(record.value())
));
}
producer.sendOffsetsToTransaction(
records.nextOffsets(),
consumer.groupMetadata()
);
producer.commitTransaction();
} catch (ProducerFencedException fatal) {
throw fatal;
} catch (KafkaException abortable) {
producer.abortTransaction();
}
}

生产代码还需处理 Rebalance、Wakeup、Fatal/Abortable Error、Backpressure 和监控。

20.10 实验验收#

应能观察:

  • 一个 Follower 故障时 ISR 缩小但仍可写;
  • ISR 低于 min ISR 后 acks=all 失败;
  • Leader 故障后从安全副本选出新 Leader;
  • Abort Transaction 对 read_committed 不可见;
  • read_uncommitted 可以读取 Abort 数据;
  • 同 Transactional ID 的旧 Producer 被 Fence;
  • 长 Open Transaction 影响 LSO。

21. 常见错误认知#

21.1 acks=all 等全部配置副本#

错误。它等待当前 ISR,并要求满足 min ISR。

21.2 RF=3 一定有三份最新数据#

错误。要看 ISR、Follower Lag 和 Under Replicated 状态。

21.3 HW 等于 Leader LEO#

错误。Leader 尾部可能尚未形成安全复制条件。

21.4 Idempotence 去重所有业务重复#

错误。它识别 Producer Retry Batch,不理解业务身份。

21.5 Transaction 会删除 Abort 数据#

错误。数据仍在 Log,通过 Marker 和 Isolation 隐藏。

21.6 read_committed 总能读到 HW#

错误。还受 LSO 限制。

21.7 Kafka EOS 保证 MySQL 只更新一次#

错误。Kafka Transaction 不覆盖 MySQL Local Transaction。


22. 本篇验收清单#

22.1 画图#

能够画出:

Producer
→ Leader Append
→ Follower Fetch
→ ISR
→ HW
→ Consumer Visibility

22.2 状态#

准确区分:

Replica Set
ISR
ELR
Leader LEO
Follower LEO
HW
LSO
Committed Offset

22.3 幂等#

能够解释:

PID
Producer Epoch
Sequence Number
Duplicate Batch
OutOfOrderSequence
Fencing

22.4 事务#

能够画出:

transactional.id
→ Transaction Coordinator
→ PID/Epoch
→ Ongoing
→ Participating Partitions
→ Prepare Commit/Abort
→ Control Marker
→ Complete

22.5 Failure Boundary#

面对“丢消息”,先问:

应用是否调用 send?
Future 是否成功?
acks 是什么?
ISR 是否满足?
是否发生 Unclean Election?
Consumer Commit 在哪里?
Retention 是否删除?
业务异常是否被吞?

面对“重复”,先问:

Kafka Log 是否重复?
应用是否重复 Send?
Producer 是否启用 Idempotence?
Consumer 是否 Replay?
数据库是否幂等?

22.6 口述验收#

不用看文章,用 45 分钟回答:

一条 acks=all Record 进入 Leader 后,Follower 如何复制,ISR 和 HW 如何限制确认与可见性;Leader 故障后 KRaft Controller 如何从 ISR、ELR 或其他候选中选举;ACK 丢失时 Producer 如何用 PID、Epoch 和 Sequence 去重;多 Partition 原子写入又如何通过 Transaction Coordinator、事务状态与 Marker 实现;最后说明 Kafka Exactly Once 为什么不能自动扩展到 MySQL 和 HTTP。


23. 全文收束:Kafka 可靠性不是一个开关#

Kafka 可靠性可以收束为五层。

Confirm Boundary#

acks + min ISR

定义 Producer 何时认为成功。

Replication Boundary#

Leader + Follower Fetch + ISR + HW

定义哪些数据可安全恢复和对 Consumer 可见。

Election Boundary#

Leader Epoch + ISR + ELR + Unclean Policy

定义故障后哪一段 Log 成为合法历史。

Retry Boundary#

PID + Producer Epoch + Sequence

让网络不确定性下的 Retry 不直接变成重复与乱序。

Atomic Visibility Boundary#

Transaction Coordinator
+ __transaction_state
+ Marker
+ LSO
+ read_committed

把多个 Kafka Partition 和 Consumer Offset 组成原子可见事务。

最终结论是:

Kafka 没有消灭失败,而是把失败转换为一组可持久化、可比较、可隔离和可恢复的状态。

但当消息离开 Kafka,进入:

MySQL
Payment
HTTP API
Agent Tool

系统仍需要:

业务幂等
Outbox / Inbox
唯一约束
状态机
补偿
对账

下一篇进入 Spring Kafka 企业工程层:

KafkaTemplate
@KafkaListener
Listener Container
AckMode
Error Handler
Blocking Retry
Retry Topic
DLT
Idempotent Consumer
Transaction Manager
Outbox / CDC

回答:

如何把 Kafka 底层可靠性机制组织成一个不会因异常重试、数据库事务和框架默认行为而失控的 Java 应用。


参考资料#

  1. Apache Kafka 4.3 — Producer Configs
  2. Apache Kafka 4.3 — KafkaProducer Javadoc
  3. Apache Kafka 4.3 — Eligible Leader Replicas
  4. Apache Kafka 4.3 — Broker Configs
  5. Apache Kafka 4.3 — Monitoring
  6. Apache Kafka 4.3 — Topic Configs
  7. Apache Kafka 4.3 — Design
  8. Apache Kafka 4.3 — Distribution Implementation
  9. Apache Kafka 4.3 — Message Format
  10. Apache Kafka Source Repository
Kafka 可靠性与事务:从副本、ISR 到 Exactly Once 的 Failure Model
https://jupiter-ws.cn/posts/backend/kafka/05_kafka_reliability_transactions/
作者
Jupiter
发布于
2026-07-18
许可协议
CC BY-NC-SA 4.0