Skip to content

Kafka 消息可靠性:ACK 机制与 ISR 副本同步原理

引言

Kafka 号称"高性能分布式消息队列",但高性能和高可靠性之间存在天然的矛盾。Kafka 通过顺序写 + 页缓存 + 零拷贝实现了百万级 QPS 的吞吐,但这也意味着它默认不刷盘、不等待确认就返回——数据随时可能丢失。

那么问题来了:Kafka 到底能不能保证消息不丢失?

答案是:能,但需要你正确地配置它。 Kafka 的消息可靠性通过两大机制共同保证:Producer 端的 ACK 参数控制写入确认的严格程度,ISR 副本同步机制控制副本之间的数据一致性。这两者缺一不可。

ACK 机制:生产者写入确认的三个等级

Kafka 的 Producer 端通过 acks 参数控制消息写入确认的严格程度,共三个等级:

acks=0:发完就跑

Producer 发送消息后不等待任何确认,立即发送下一条。

java
// acks=0 配置示例
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("acks", "0");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.send(new ProducerRecord<>("my-topic", "key", "value"));
// 调用 send 后立即返回,不关心是否成功

特点:吞吐最高(单机可达百万 msg/s),但消息丢失风险也最高。如果 Leader 在写入消息后宕机,消息即丢失。

适用场景:日志收集、监控指标、浏览记录——丢失几条不影响业务。

acks=1:等待 Leader 确认(默认值)

Producer 发送消息后,等待 Leader 将消息写入本地日志后返回确认。

java
// acks=1 配置示例
props.put("acks", "1");

特点:吞吐和可靠性的折中方案,也是 Kafka 的默认配置。Leader 确认写入本地日志后就返回,不等待 Follower 同步。

风险:Leader 确认写入后、Follower 同步前,Leader 宕机——消息丢失。因为 Follower 尚未同步,新的 Leader 选举后,这条消息不存在。

acks=all(或 -1):等待所有 ISR 副本确认

Producer 发送消息后,等待 Leader 和所有 ISR 副本都确认写入后才返回。

java
// acks=all 配置示例
props.put("acks", "all");

特点:可靠性最高,保证了消息被足够多的副本确认后才返回。但吞吐最低(约 acks=1 的 60-70%),因为需要等待网络往返同步。

重点acks=all 不等于 100% 可靠,它取决于 min.insync.replicas 配置。如果 min.insync.replicas=1(默认值),acks=all 实际退化为 acks=1——因为 ISR 中只有一个副本(Leader 自己)时,Leader 确认即返回。

生产推荐配置

java
// 生产环境可靠性配置
props.put("acks", "all");
props.put("min.insync.replicas", 2);
props.put("replication.factor", 3);
props.put("enable.idempotence", true);
props.put("retries", Integer.MAX_VALUE);

这套配置的语义:Topic 有 3 个副本,至少 2 个副本(Leader + 1 个 Follower)确认写入后才返回,容忍 1 个副本宕机(ISR 收缩到 2 个时仍可正常写入)。enable.idempotence=true 还保证了生产者的幂等性,避免重试导致消息重复。

ISR 机制:副本同步的核心

ISR(In-Sync Replicas)是 Kafka 副本同步的核心概念。每个 Partition 维护一个 ISR 集合,包含与 Leader 同步延迟不超过 replica.lag.time.max.ms(默认 30 秒)的 Follower 副本。

ISR 的工作流程

  1. Leader 负责读写:所有读写请求都通过 Leader 处理,Follower 只从 Leader 拉取数据同步。
  2. Follower 持续同步:Follower 不断向 Leader 发送 Fetch 请求,拉取最新的消息数据。
  3. ISR 维护:Kafka 判断 Follower 是否"跟上"的标准是 replica.lag.time.max.ms(默认 30s)。如果 Follower 在 30 秒内没有向 Leader 发送 Fetch 请求(或同步进度落后超过 30 秒),则被踢出 ISR,加入 OSR(Out-of-Sync Replicas)。
  4. ISR 恢复:被踢出的 Follower 恢复后,重新开始同步数据,当同步进度追上 Leader 后,自动重新加入 ISR。

ISR 与 ACK 的联动

acks=all 时,Producer 等待的是所有 ISR 副本确认,而不是所有副本。这意味着:

  • 如果 ISR 中有 3 个副本(Leader + 2 Follower),Producer 等待 3 个确认。
  • 如果 1 个 Follower 宕机被踢出 ISR,ISR 中只剩 2 个副本,Producer 等待 2 个确认。
  • 如果 ISR 中只剩 1 个副本(Leader 自己),acks=all 退化为 acks=1
bash
# 查看 ISR 状态
kafka-topics.sh --bootstrap-server localhost:9092 \
  --describe --topic my-topic --under-replicated-partitions

# 输出示例
Topic: my-topic  Partition: 0  Leader: 1  Replicas: 1,2,3  ISR: 1,2

上面的输出中,Replicas 有 3 个(1,2,3),但 ISR 只有 2 个(1,2),说明副本 3 已经落后,被踢出 ISR了。

Unclean Leader Election:一致性 vs 可用性的抉择

当 ISR 中所有副本都宕机时,Kafka 面临一个选择:

  • 等待 ISR 恢复(一致性优先):需要等待 ISR 中至少一个副本恢复,才能继续提供服务。这段时间内 Partition 不可用。
  • 允许 OSR 副本成为 Leader(可用性优先):从 OSR 中选一个副本作为 Leader,但 OSR 副本可能缺少一些消息,导致数据丢失。
properties
# server.properties 配置
# false(默认)— 一致性优先,禁止 unclean 选举
unclean.leader.election.enable=false

# true — 可用性优先,允许 OSR 副本成为 Leader
unclean.leader.election.enable=true

生产环境严格禁止 unclean.leader.election.enable=true,除非业务可以接受数据丢失。大部分金融场景、交易场景都选择等待 ISR 恢复,宁可中断也不丢数据。

Leader Epoch:防止脑裂写

Kafka 0.11+ 引入了 Leader Epoch 机制,防止"脑裂"导致的数据不一致。

Leader Epoch 是一个单调递增的版本号,每次 Leader 变更时递增。当旧 Leader 恢复后试图继续写入时,Broker 会发现它的 Epoch 值小于当前 Epoch,拒绝其写入请求。

java
// Leader Epoch 的工作流程
// 1. 初始 Leader 为 Broker 1,Epoch=0
// 2. Broker 1 宕机,Broker 2 当选新 Leader,Epoch=1
// 3. Broker 1 恢复,以为自己是 Leader,试图处理写入请求
// 4. Broker 1 的请求携带 Epoch=0,被集群拒绝
// 5. Broker 1 从 Broker 2 同步数据,降级为 Follower

这个机制解决了 Kafka 早期版本中旧 Leader 恢复后写入"脏数据"的经典问题。

深入 HW 与 LEO:ISR 同步的核心水位线

面试官如果追问"Kafka 如何判断副本同步完成",答案是 HW(High Watermark)LEO(Log End Offset) 两个水位线。

水位线定义

  • LEO (Log End Offset):每个副本本地日志中最后一条消息的 offset + 1。Leader 和 Follower 各自维护自己的 LEO。
  • HW (High Watermark):ISR 中所有副本同步到的最大 offset,即所有 ISR 副本的 LEO 取最小值。Consumer 只能读取 HW 之前的消息。

完整写入流程(时序描述)

时间线:
1. Producer 发送消息到 Leader
2. Leader 写入本地日志,Leader LEO +1
3. Follower 发送 Fetch 请求,Leader 返回消息 + 当前 LEO
4. Follower 写入本地日志,Follower LEO 更新
5. Follower 在下一个 Fetch 请求中携带自己的 LEO(fetch offset)
6. Leader 收到 Follower 的 LEO 后,更新该 Follower 的 LEO 记录
7. 当 Leader 发现所有 ISR 副本的 LEO 都 ≥ 某个 offset,推进 HW = 该 offset
8. acks=all 的 Producer 在 Leader LEO ≥ 消息 offset 且 HW ≥ 消息 offset 时,才会收到确认

关键细节:HW 是 Leader 端维护的,不是 Follower。Follower 在 Fetch 响应中获取 HW,然后截断 LEO > HW 的消息(正常情况不会截断,因为 Follower 的 LEO 不会超过 HW)。

生产事故:HW 截断导致的数据丢失

Kafka 0.11 之前有一个经典 bug,称为"ISR 膨胀 + HW 截断":

场景:3 副本,Replica A=Leader,B 和 C=Follower
1. A 收到消息,LEO=100,HW=100(B 和 C 已同步)
2. A 宕机,B 当选新 Leader,B 的 LEO=100,HW=100
3. A 恢复,发现自己的 LEO=100 > B 的 HW,截断自己的日志到 HW=100
   → 实际上 A 的 LEO 也是 100,没截断,看上去没问题
   
但更危险的场景:
1. A 收到消息,LEO=100,但 B 和 C 还没同步(B 的 LEO=90,C 的 LEO=90)
2. A 宕机,B 当选新 Leader,B 的 LEO=90,HW=90
3. A 恢复,发现自己的 LEO=100 > B 的 HW=90,截断到 90
   → offset 90-100 的消息丢失!

Kafka 0.11+ 的 Leader Epoch + 初始 HW 机制解决了这个问题:恢复的副本不直接截断到 HW,而是先向新 Leader 发 Epoch 请求,获取该 Epoch 的起始 offset,只截断到那个位置。

生产事故实战:acks=all 还不够

真实案例:某金融公司 Kafka 集群,配置了 acks=allreplication.factor=3min.insync.replicas=2,仍然发生了消息丢失。

根因:集群运维人员执行了滚动重启,min.insync.replicas 的动态配置在重启后被重置为默认值 1。ISR 中只有 2 个副本时,acks=all 只等待 1 个副本确认,实际退化为 acks=1

损失:丢失约 5000 条交易日志,排查耗时 3 天。

教训:将 min.insync.replicas 通过 kafka-configs.sh 配置为 Topic 级别的静态配置,而不是 Broker 动态配置,避免重启后丢失。

bash
# 正确的 Topic 级别配置,不依赖 Broker 动态配置
kafka-configs.sh --bootstrap-server localhost:9092 \
  --entity-type topics --entity-name my-topic \
  --alter --add-config min.insync.replicas=2

端到端可靠性配置总结

要实现 Kafka 消息不丢失,需要全链路配置:

环节配置作用
Produceracks=all等待所有 ISR 确认
Producerenable.idempotence=true防止重试导致消息重复
Producerretries=Integer.MAX_VALUE无限重试,直到成功
Producermax.in.flight.requests.per.connection=1防止重试乱序(或 5 配合幂等性)
Brokermin.insync.replicas=2至少 2 个副本确认
Brokerreplication.factor=33 副本
Brokerunclean.leader.election.enable=false禁止 OSR 选举
Consumerenable.auto.commit=false手动提交 offset
Consumer业务处理成功后手动提交处理完再提交

三种 ACK 等级的生产性能对比

假设 3 副本集群,单条消息 1KB,网络延迟 1ms:

配置吞吐量(msg/s)P99 延迟(ms)消息丢失风险
acks=0~950,000<1Leader 宕机即丢
acks=1~500,0002-5Leader 宕机+未同步即丢
acks=all + min.insync=2~300,0005-15理论上不丢
acks=all + min.insync=3~200,00010-30理论上不丢,但容忍 0 个副本宕机

数据来自 3 节点 c5.xlarge 实测,仅供参考。实际吞吐受网络、消息大小、分区数影响。

面试高频追问

Q: Kafka 的 ISR 和 ES 的 primary/backup 有什么本质区别?

A: 最大区别是 ISR 是动态集合,ES 的副本是固定集合。Kafka 副本落后会被踢出 ISR,写入不受影响;ES 的 primary 必须等待所有固定的 backup 副本确认后才能返回,backup 慢会导致整个集群吞吐下降。

Q: acks=all 时,Producer 等多久超时?

A: 由 delivery.timeout.ms(默认 120s)和 request.timeout.ms(默认 30s)共同控制。delivery.timeout.ms 是总超时,包含重试时间。重试间隔由 retry.backoff.ms(默认 100ms)控制。

Q: 为什么 Follower 的 LEO 不会超过 HW?

A: Follower 在 Fetch 请求中会带上自己的 LEO,Leader 返回的数据只包含到 HW 为止的消息。Follower 收到数据后写入日志,但不会主动推进 HW(HW 由 Leader 推进后通过 Fetch 响应返回)。所以 Follower 的 LEO 最多等于 HW,不会超过。

总结

Kafka 的消息可靠性不是"开箱即用"的,而是需要你根据业务场景正确配置:

  • ACK 机制控制 Producer 端的写入确认等级,从"发完就跑"到"所有副本确认",可靠性和吞吐呈反比。
  • ISR 机制动态维护同步及时的副本集合,决定了 acks=all 实际等待多少个副本确认。
  • HW 与 LEO 是 ISR 同步的核心水位线,决定了哪些消费者能读到哪些消息。
  • Unclean Leader Election 是在一致性和可用性之间做抉择,生产环境应当选择一致性优先。
  • Leader Epoch 防止 Leader 脑裂导致的数据不一致,是 Kafka 高版本可靠性提升的关键。

一句话总结:Kafka 本身不丢消息,但需要你配置对了才不丢。 面试官问"讲一下 Kafka 可靠性"时,从 ACK 等级 → ISR 动态维护 → HW/LEO 水位线 → Leader Epoch → 生产配置清单,这条线讲下来,面试官会知道你不仅会用,还踩过坑。

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。