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 重复调用 sendRequest 没有到达 BrokerLeader 已追加但 ACK 丢失Follower 落后或掉出 ISRLeader 在返回成功后永久故障旧 Leader 在网络分区后恢复Consumer 读到未完成事务的数据业务数据库成功,但 Offset Commit 失败完成本篇后,你应该能够:
- 从 Failure Boundary 分析消息丢失、重复与乱序;
- 准确解释
acks=0、acks=1、acks=all; - 说明 Replication Factor、Leader、Follower 和 Replica Fetch 的关系;
- 区分 Replica Set、ISR、ELR、HW、LEO 和 LSO;
- 解释
min.insync.replicas与acks=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=all、retries>0 且 max.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 已追加到 LogProduceResponse ↓ ACK 在网络中丢失Producer Timeout
Producer 无法区分:
Request 根本没有到达 BrokerBroker 收到但尚未追加Broker 已追加,ACK 在回程丢失Broker 已追加并复制,但连接中断如果不重试,第一种情况会丢消息;如果重试,第三种情况可能重复。
这说明:
分布式发送最难的不是服务器明确返回失败,而是客户端无法判断服务器执行到了哪一步。
Idempotent Producer 没有让网络变可靠。它解决的是:
当 Producer 因不确定结果重试同一个 Batch 时,Broker 能否识别它已经追加过。
1. 一条消息的 Failure Boundary

“Kafka 会不会丢消息”不是一个足够精确的问题。至少要拆成七段。
1.1 Application Boundary
应用可能因为 HTTP 重试、任务重放或代码错误重复调用:
producer.send(event);Kafka Idempotence 无法理解两个不同 send() 在业务上是不是同一事件。业务层仍需要:
Event IDRequest ID订单号 + 事件类型业务状态机幂等记录1.2 Producer Boundary
Record 可能仍处于:
序列化阶段RecordAccumulator等待 BufferIn-flight RequestRetry BackoffDelivery Timeoutsend() 返回 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 WatermarkIsolation LevelLast Stable OffsetPosition限制。
1.7 External Side Effect Boundary
消费后可能:
更新 MySQL调用支付 API修改库存执行 Agent Tool这些副作用不属于 Kafka Log,Kafka Transaction 不能自动覆盖它们。
因此排障时固定问:
失败发生在哪一段?这一段的 Source of Truth 是什么?调用方能否确定执行结果?Retry 是否会产生重复?恢复状态由谁保存?2. acks:Producer 何时认为写入完成

acks 定义 Producer 等待的 Broker 确认条件。它不是可靠性的唯一开关,可靠性还取决于:
Replication FactorISRmin.insync.replicasLeader ElectionRetryIdempotence故障域分布应用是否检查 Future2.1 acks=0
acks=0Producer 不等待 Broker Response。消息进入 Socket Buffer 后就认为已发送。
结果:
Broker 是否收到:未知是否被拒绝:未知Retry 无法依赖 Broker 错误RecordMetadata.offset=-1只适合可容忍丢失的低价值遥测。
2.2 acks=1
acks=1Leader 本地追加后立即返回,不等待 Follower。
Leader 追加 A→ 返回成功→ Follower 尚未复制→ Leader 永久故障→ 新 Leader 没有 AProducer 已经看到成功,但 Record 可能丢失。
2.3 acks=all
acks=all等价于 -1。Leader 等待当前 ISR 全部确认,并检查 min.insync.replicas。
Kafka 4.3 Producer 默认使用 acks=all。
2.4 acks=all 不是固定等待两个副本
假设:
RF=3ISR={1,2,3}min ISR=2acks=all 等待 1、2、3,而不是只等两个。
如果:
ISR={1,2}等待 1、2。
如果:
ISR={1}因为不满足 min ISR=2,写入被拒绝。
2.5 成功响应不等于每块磁盘都同步 fsync
需要区分:
Producer Response SuccessConsumer VisibilityPhysical Device FlushKafka 的核心可靠性依赖副本协议、确认条件和选举约束,不是每条消息都同步刷所有物理磁盘。
3. Kafka 副本模型:Leader 与 Follower
一个 Partition 的:
Replication Factor=N表示它配置了 N 个 Replica,其中一个 Leader,其余 Follower。

3.1 Producer 只写 Leader
所有写入经 Leader 排序并分配 Offset,保证 Partition 只有一个合法追加顺序。
3.2 Follower 主动 Fetch
Follower LEO ↓ FetchRequest(offset=LEO)Leader ↓ 返回 RecordBatchFollower 追加本地 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。

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 Replicas5. High Watermark、LEO 与可见性
5.1 Log End Offset
每个 Replica 都有自己的 LEO:
Leader LEO=12Follower A LEO=10Follower 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 EpochISR 变化严格 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:持久性与可用性的选择

假设:
RF=3min ISR=26.1 ISR={1,2,3}
满足写入条件,acks=all 等待三个当前 ISR。
6.2 ISR={1,2}
仍满足 min ISR=2,继续写入,等待两个 ISR。
6.3 ISR={1}
写入被拒绝,Producer 可能收到:
NotEnoughReplicasExceptionNotEnoughReplicasAfterAppendException6.4 为什么拒绝写入
只剩 Leader 仍继续确认:
唯一 Leader 写入→ Producer 成功→ Leader 永久丢失→ 无副本可恢复拒绝写入是在明确表达:
当前集群无法继续满足声明的持久性级别。
6.5 权衡
提高 min ISR:
+ 已确认数据更难丢失- 故障时更早停止写入降低:
+ 故障期间更容易保持可写- 可安全恢复副本更少常见关键业务组合:
RF=3min ISR=2acks=allunclean election disabledRack Awareness7. Leader Election:ISR、ELR 与 Unclean

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 Leader7.3 ELR 与 Unclean 的根本差异
ELR→ 虽不在 ISR,但仍被证明包含安全提交历史
Unclean→ 可能缺失已确认数据ELR 目标是在安全前提下改善可用性;Unclean 用数据风险换可用性。
7.4 Unclean Leader Election
unclean.leader.election.enable=true假设:
旧 Leader: 0..100Follower: 0..80Follower 成为新 Leader 后,81~100 会从合法历史消失。
关键交易通常不应开启 Unclean Election。
7.5 Leader Epoch 与旧 Leader Fencing
网络分区后,旧 Leader 可能仍认为自己是 Leader。KRaft Controller 通过 Leader Epoch 和 Metadata 隔离旧 Leader。
旧 Broker 恢复后必须:
确认新 Epoch截断不合法尾部作为 Follower 重新追赶8. Retry、重复与乱序

8.1 ACK 丢失导致重复
A 已追加ACK 丢失Producer Retry A没有去重状态时,Broker 无法判断它是新 Batch 还是旧 Batch 重试。
8.2 多 In-flight 导致乱序
同一 Partition:
A 失败等待重试B 后发先成功A 重试成功可能写成:
B → AKafka 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

Kafka 4.3 在没有冲突配置时默认:
enable.idempotence=true依赖:
acks=allretries>0max.in.flight.requests.per.connection<=59.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 / Fenced9.5 去重作用域
Idempotence 能解决:
同一 Producer Session同一 TopicPartition同一 Batch 的 Broker Retry不能解决:
两个 Producer 发送相同 Event应用重启后重新构造事件Consumer 重复更新数据库人工 ReplayHTTP 上游重复请求9.6 为什么要求 acks=all
如果使用 acks=1,Leader 确认 Sequence 后永久故障,而 Follower 未复制 Batch 和 Producer State,新 Leader 无法可靠延续去重状态。
9.7 Broker 状态恢复
源码重点:
ProducerStateManagerProducerAppendInfoBatchMetadataVerificationStateEntry负责 Sequence 校验、Epoch、Duplicate Detection、Transaction 状态和 Snapshot 恢复。
10. Idempotence 为什么不等于 Transaction

Idempotence 解决单 Partition Retry 去重。
业务可能需要:
P0: OrderCreatedP1: InventoryReserved如果 P0 成功、P1 失败,即使两个写入各自幂等,业务仍是部分成功。
Kafka Transaction 提供:
多个 Topic多个 PartitionConsumer Offset的原子可见性。
10.1 事务不是磁盘回滚
事务 Record 会先追加进普通 Log。Abort 时不物理删除,而是使用:
Transaction MetadataAbort MarkerConsumer Isolation让 read_committed Consumer 跳过。
Kafka 事务控制的是原子可见性,不是 Undo 写入。
11. Transactional Producer API
配置:
transactional.id=fulfillment-worker-01设置后自动启用 Idempotence。
11.1 initTransactions()
producer.initTransactions();它会:
- 查找 Transaction Coordinator;
- 完成或 Abort 同 ID 的旧事务;
- 获取 PID 与 Epoch;
- 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 与状态机

12.1 Coordinator 定位
hash(transactional.id)→ __transaction_state Partition→ 该 Partition Leader→ Transaction CoordinatorKafka 4.3 默认:
transaction.state.log.num.partitions=50transaction.state.log.replication.factor=3transaction.state.log.min.isr=212.2 Transaction Metadata
概念上包括:
transactional.idPIDProducer EpochTransaction TimeoutState参与 TopicPartitionsStart Timestamp12.3 状态机
Empty→ Ongoing ├→ PrepareCommit → CompleteCommit └→ PrepareAbort → CompleteAbort完成过程:
- 持久化即将 Commit/Abort 的决定;
- 向参与 Partition 写 Control Marker;
- 所有 Marker 完成后写最终状态。
12.4 AddPartitionsToTxn
Producer 首次向某 Partition 写事务数据时,需要将其加入当前事务参与集合,因为 Coordinator 最终必须知道向哪些 Partition 写 Marker。
12.5 Transaction Timeout
Broker 默认上限:
transaction.max.timeout.ms=900000长事务会:
阻挡 LSO延迟 read_committed增加恢复成本占用 Transaction StateKafka Transaction 不适合包裹分钟级人工审批或长时间 Agent Workflow。
13. Commit/Abort Marker、LSO 与隔离读取

13.1 事务数据先追加
事务 Record 已占用正常 Offset。
13.2 Control Marker
Commit 时写:
COMMIT MarkerAbort 时写:
ABORT MarkerControl 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

不用事务时:
输出成功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 的边界

“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 Publisher15.5 Kafka → HTTP / Payment
需要下游支持:
Idempotency-KeyRequest ID状态查询业务状态机补偿与对账15.6 Effectively Once
跨系统无法建立全局事务时,常见目标:
At-Least-Once Delivery+ Idempotent Side Effect= Effectively Once16. 丢失、重复与乱序的来源
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 进入不同 PartitionTopic 扩容改变 Key 映射禁用 Idempotence + 多 In-flight Retry同 Partition 并发业务 WorkerRetry Topic / DLT多 Topic 合流事件时间与到达时间不同17. 可靠性配置矩阵
17.1 Producer
acks=allenable.idempotence=trueretries=2147483647max.in.flight.requests.per.connection=5delivery.timeout.ms=120000request.timeout.ms=30000多数场景保留 Kafka 4.3 默认值,主要通过 SLO 调整 delivery.timeout.ms,不要盲目设置巨大 Retry 次数和应用级重发。
17.2 Topic
关键事件常见:
RF=3min.insync.replicas=2unclean.leader.election.enable=falseRack Awareness17.3 Transaction Internal Topic
生产默认:
transaction.state.log.replication.factor=3transaction.state.log.min.isr=2单节点实验临时降低后不要复制到生产。
17.4 Consumer
事务下游:
isolation.level=read_committedenable.auto.commit=false具体 Offset 是否手动提交取决于 Kafka Transaction、Streams 或框架 Container Transaction。
18. Metrics 与告警

18.1 Broker 副本
UnderReplicatedPartitionsUnderMinIsrPartitionCountOfflineReplicaCountIsrShrinksPerSecIsrExpandsPerSecFailedIsrUpdatesPerSecReplicaFetcher MaxLag正常情况下 Offline Replica 为 0,非故障恢复期 ISR Shrink/Expansion 应接近 0。
18.2 Producer
record-error-raterecord-retry-raterequest-latencyrequests-in-flightrecord-queue-timebufferpool-wait-timeproduce-throttle-time18.3 Transaction
transaction-start-ratetransaction-commit-ratetransaction-abort-ratetransaction-durationOpen TransactionLSO LagTransaction Coordinator Load实际指标名应以具体 Kafka 4.3.x 客户端和 Reporter 输出为准。
18.4 告警
至少包括:
UnderMinISR > 0OfflineReplicaCount > 0ISR Shrink 突增Unclean ElectionProducer Error/Retry 激增Transaction Abort 激增Transaction Duration 超限read_committed LSO Lag19. 源码阅读地图
固定 Kafka 4.3.x Tag 阅读,避免 trunk 方法变化。
19.1 Produce 与 Replica
KafkaApisReplicaManagerPartitionUnifiedLog关注:
- ProduceRequest 如何进入 Append;
acks如何影响 Delayed Produce;- min ISR 在 Append 前后如何校验;
- ProduceResponse 何时完成;
- Leader Epoch 如何传递。
19.2 Replica Fetch
ReplicaFetcherManagerReplicaFetcherThreadAbstractFetcherThread关注:
- Follower Fetch Offset;
- FetchResponse 追加;
- Diverging Epoch;
- Log Truncation;
- ISR 恢复。
19.3 ISR、HW 与 ELR
PartitionReplicaManagerAlterPartitionManagerKRaft Partition Registration关注 ISR Expand/Shrink、HW 更新、Controller 确认和 ELR 候选状态。
19.4 Idempotent Producer
Client:
TransactionManagerProducerIdAndEpochSenderProducerBatchBroker:
ProducerStateManagerProducerAppendInfoBatchMetadataVerificationStateEntry19.5 Transaction Coordinator
TransactionCoordinatorServiceTransactionCoordinatorShardTransactionMetadataTransactionStateTransactionLogTransactionMarkerChannelManager关注:
transactional.id 分区InitProducerIdAddPartitionsPrepare Commit/AbortMarkerTimeout AbortCoordinator Failover19.6 Isolation Fetch
UnifiedLog.readFetchIsolationAbortedTxnTransactionIndexCompletedTxn推荐顺序:
acks/min ISR 文档→ Partition Append→ Replica Fetcher→ ISR/HW/ELR→ ProducerStateManager→ Client TransactionManager→ Transaction Coordinator→ Marker/LSO/Isolation Fetch20. 可运行实验:三 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:启动:
docker compose up -d20.2 创建 Topic
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查看:
docker exec kafka-r1 \ /opt/kafka/bin/kafka-topics.sh \ --bootstrap-server kafka-1:9092 \ --describe \ --topic reliable-events20.3 停止一个 Follower
使用 acks=all 持续生产,然后停止一个非 Leader Broker:
docker stop kafka-r3观察 ISR 缩小。RF=3、ISR=2、min ISR=2 时仍可写。
20.4 只剩一个 ISR
再停止另一个 ISR Broker。如果 Leader 仍在线但 ISR 只剩一个,acks=all 应被拒绝。
恢复:
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 对比事务可见性
分别启动:
--consumer-property isolation.level=read_uncommitted和:
--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 Visibility22.2 状态
准确区分:
Replica SetISRELRLeader LEOFollower LEOHWLSOCommitted Offset22.3 幂等
能够解释:
PIDProducer EpochSequence NumberDuplicate BatchOutOfOrderSequenceFencing22.4 事务
能够画出:
transactional.id→ Transaction Coordinator→ PID/Epoch→ Ongoing→ Participating Partitions→ Prepare Commit/Abort→ Control Marker→ Complete22.5 Failure Boundary
面对“丢消息”,先问:
应用是否调用 send?Future 是否成功?acks 是什么?ISR 是否满足?是否发生 Unclean Election?Consumer Commit 在哪里?Retention 是否删除?业务异常是否被吞?面对“重复”,先问:
Kafka Log 是否重复?应用是否重复 Send?Producer 是否启用 Idempotence?Consumer 是否 Replay?数据库是否幂等?22.6 口述验收
不用看文章,用 45 分钟回答:
一条
acks=allRecord 进入 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,进入:
MySQLPaymentHTTP APIAgent Tool系统仍需要:
业务幂等Outbox / Inbox唯一约束状态机补偿对账下一篇进入 Spring Kafka 企业工程层:
KafkaTemplate@KafkaListenerListener ContainerAckModeError HandlerBlocking RetryRetry TopicDLTIdempotent ConsumerTransaction ManagerOutbox / CDC回答:
如何把 Kafka 底层可靠性机制组织成一个不会因异常重试、数据库事务和框架默认行为而失控的 Java 应用。
参考资料
- Apache Kafka 4.3 — Producer Configs
- Apache Kafka 4.3 — KafkaProducer Javadoc
- Apache Kafka 4.3 — Eligible Leader Replicas
- Apache Kafka 4.3 — Broker Configs
- Apache Kafka 4.3 — Monitoring
- Apache Kafka 4.3 — Topic Configs
- Apache Kafka 4.3 — Design
- Apache Kafka 4.3 — Distribution Implementation
- Apache Kafka 4.3 — Message Format
- Apache Kafka Source Repository