Kafka 消息可靠性:ACK 机制与 ISR 副本同步原理
引言
Kafka 号称"高性能分布式消息队列",但高性能和高可靠性之间存在天然的矛盾。Kafka 通过顺序写 + 页缓存 + 零拷贝实现了百万级 QPS 的吞吐,但这也意味着它默认不刷盘、不等待确认就返回——数据随时可能丢失。
那么问题来了:Kafka 到底能不能保证消息不丢失?
答案是:能,但需要你正确地配置它。 Kafka 的消息可靠性通过两大机制共同保证:Producer 端的 ACK 参数控制写入确认的严格程度,ISR 副本同步机制控制副本之间的数据一致性。这两者缺一不可。
ACK 机制:生产者写入确认的三个等级
Kafka 的 Producer 端通过 acks 参数控制消息写入确认的严格程度,共三个等级:
acks=0:发完就跑
Producer 发送消息后不等待任何确认,立即发送下一条。
// 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 将消息写入本地日志后返回确认。
// acks=1 配置示例
props.put("acks", "1");特点:吞吐和可靠性的折中方案,也是 Kafka 的默认配置。Leader 确认写入本地日志后就返回,不等待 Follower 同步。
风险:Leader 确认写入后、Follower 同步前,Leader 宕机——消息丢失。因为 Follower 尚未同步,新的 Leader 选举后,这条消息不存在。
acks=all(或 -1):等待所有 ISR 副本确认
Producer 发送消息后,等待 Leader 和所有 ISR 副本都确认写入后才返回。
// 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 确认即返回。
生产推荐配置
// 生产环境可靠性配置
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 的工作流程
- Leader 负责读写:所有读写请求都通过 Leader 处理,Follower 只从 Leader 拉取数据同步。
- Follower 持续同步:Follower 不断向 Leader 发送 Fetch 请求,拉取最新的消息数据。
- ISR 维护:Kafka 判断 Follower 是否"跟上"的标准是
replica.lag.time.max.ms(默认 30s)。如果 Follower 在 30 秒内没有向 Leader 发送 Fetch 请求(或同步进度落后超过 30 秒),则被踢出 ISR,加入 OSR(Out-of-Sync Replicas)。 - 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。
# 查看 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 副本可能缺少一些消息,导致数据丢失。
# 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,拒绝其写入请求。
// 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=all、replication.factor=3、min.insync.replicas=2,仍然发生了消息丢失。
根因:集群运维人员执行了滚动重启,min.insync.replicas 的动态配置在重启后被重置为默认值 1。ISR 中只有 2 个副本时,acks=all 只等待 1 个副本确认,实际退化为 acks=1。
损失:丢失约 5000 条交易日志,排查耗时 3 天。
教训:将 min.insync.replicas 通过 kafka-configs.sh 配置为 Topic 级别的静态配置,而不是 Broker 动态配置,避免重启后丢失。
# 正确的 Topic 级别配置,不依赖 Broker 动态配置
kafka-configs.sh --bootstrap-server localhost:9092 \
--entity-type topics --entity-name my-topic \
--alter --add-config min.insync.replicas=2端到端可靠性配置总结
要实现 Kafka 消息不丢失,需要全链路配置:
| 环节 | 配置 | 作用 |
|---|---|---|
| Producer | acks=all | 等待所有 ISR 确认 |
| Producer | enable.idempotence=true | 防止重试导致消息重复 |
| Producer | retries=Integer.MAX_VALUE | 无限重试,直到成功 |
| Producer | max.in.flight.requests.per.connection=1 | 防止重试乱序(或 5 配合幂等性) |
| Broker | min.insync.replicas=2 | 至少 2 个副本确认 |
| Broker | replication.factor=3 | 3 副本 |
| Broker | unclean.leader.election.enable=false | 禁止 OSR 选举 |
| Consumer | enable.auto.commit=false | 手动提交 offset |
| Consumer | 业务处理成功后手动提交 | 处理完再提交 |
三种 ACK 等级的生产性能对比
假设 3 副本集群,单条消息 1KB,网络延迟 1ms:
| 配置 | 吞吐量(msg/s) | P99 延迟(ms) | 消息丢失风险 |
|---|---|---|---|
| acks=0 | ~950,000 | <1 | Leader 宕机即丢 |
| acks=1 | ~500,000 | 2-5 | Leader 宕机+未同步即丢 |
| acks=all + min.insync=2 | ~300,000 | 5-15 | 理论上不丢 |
| acks=all + min.insync=3 | ~200,000 | 10-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 → 生产配置清单,这条线讲下来,面试官会知道你不仅会用,还踩过坑。