大家好,我是小富。
《十万个why》系列持续更新中
上周五晚上,订单组的同事被一个告警炸醒了。线上有个 Kafka Topic 配了 3 副本、acks=all,生产者那边发送回调全是 success,监控上看 ISR 列表也一直是 3 个副本。结果半夜 Broker-2 整机宕机,重启之后拉出日志一对比,发现有将近 1200 条订单消息在消费端根本找不到。
排查到最后,原因不是网络丢包,也不是消费端漏消费,而是那 1200 条消息压根就没真正持久化到三个副本上。
acks=all 不是万能保险箱,它只是 Kafka 给你的一个承诺框架,框架里的每一个零件你都得自己拧紧。
先搞清楚 acks=all 到底承诺了什么
很多人对 acks=all 的理解是:生产者发一条消息,三个副本都写入磁盘了才返回成功。
这个理解少了两个关键前提。
acks=all 的真实含义是:Leader 收到消息后,等待 ISR(In-Sync Replicas)列表中的所有副本都确认收到了这条消息,才给生产者返回 ack。
注意两个细节:
- "所有副本" 指的是 ISR 列表里的副本,不是配置的总副本数。如果 ISR 里只剩 Leader 自己,那 all 就等于 1。
- "确认收到" 是写进了操作系统的 Page Cache,不是刷盘到磁盘。
这两个细节,就是消息丢失的根源。
第一个坑:min.insync.replicas 没设对
这是最常见的配置失误。
Kafka 有个参数叫 min.insync.replicas,默认值是 1。它的含义是:当 acks=all 时,ISR 列表里至少要有多少个副本确认,消息才算写入成功。
默认值是 1 意味着什么?只要 Leader 自己确认了就算 all。
来看一下整个写入流程:
Producer Broker-0(Leader) Broker-1(Follower) Broker-2(Follower)
| | | |
|--- send(msg, acks=all) ---->| | |
| |-- 写入 Page Cache | |
| | | |
| | (min.insync.replicas=1, ISR 中 Leader 已确认) |
| | | |
|<---- ack: success ---------| | |
| | | |
| [Follower 还没来得及拉取这条消息] | |
| | | |
| [Broker-0 宕机] | |
| X | |
| |--- 被选为新 Leader |
| | (没有这条消息) |
看到了吧,acks=all + 3 副本,但因为 min.insync.replicas=1,消息只写到了 Leader 就返回成功了。Leader 一宕机,消息就没了。
生产者那边的回调日志显示 success,监控面板上一切正常,但数据已经丢了。这就是为什么只看 acks=all 是不够的。
正确的做法是把 min.insync.replicas 设成 2:
# broker 级别或 topic 级别配置
min.insync.replicas=2
这样配了之后,acks=all 要求 ISR 中至少 2 个副本确认才返回成功。3 副本的情况下,Leader + 至少 1 个 Follower 都确认了才算写入成功。
但这里有个连带影响你得注意:如果 ISR 里只剩下 1 个副本(比如两个 Follower 都挂了),生产者会直接收到 NotEnoughReplicasException,消息写入失败。
// 生产者发送消息时可能收到的异常
org.apache.kafka.common.errors.NotEnoughReplicasException:
Messages are rejected since there are fewer in-sync replicas than required.
有些团队觉得这样不好,写入失败影响业务,于是把 min.insync.replicas 调回 1。说白了就是拿数据可靠性换可用性,你得自己权衡。
第二个坑:ISR 缩水你可能完全不知道
就算你把 min.insync.replicas 设成了 2,ISR 列表也不是一成不变的。
Follower 副本需要持续从 Leader 拉取数据来保持同步。如果某个 Follower 因为网络抖动、GC 停顿、磁盘 I/O 太高等原因,拉取速度跟不上 Leader,Kafka 会把它从 ISR 列表里踢出去。
控制这个行为的参数是 replica.lag.time.max.ms,默认 30 秒。Follower 超过 30 秒没有追上 Leader 的最新 offset,就会被踢出 ISR。
问题是,ISR 从 3 缩到 2 的时候,生产者端没有任何感知。日志不会报错,回调还是 success,min.insync.replicas=2 的约束依然满足。
但这时候你的冗余已经从"可以容忍 1 台宕机"变成了"任何一台再出问题就完蛋"。如果在这个窗口期里 Leader 也宕机了,消息就丢了。
所以光配参数不够,还得加监控。写一个定时检查脚本或者接入 Kafka 的 JMX 指标,盯住每个 Topic 的 ISR 数量:
# 查看 Topic 的分区详情,重点关注 Isr 列表
kafka-topics.sh --describe --topic order-topic \
--bootstrap-server 192.168.1.10:9092
# 输出示例:
# Topic: order-topic Partition: 0 Leader: 0 Replicas: 0,1,2 Isr: 0,1
# ^^^^^^^^
# Isr 只剩 2 个了,Broker-2 掉队了
ISR 数量低于副本数的时候就应该告警了,别等到宕机才发现。
第三个坑:unclean.leader.election 允许落后副本当 Leader
这个参数叫 unclean.leader.election.enable,在 Kafka 0.11 之前默认是 true,0.11 之后默认改成了 false。但很多公司的集群是从老版本升上来的,这个参数可能还是 true。
它的含义是:当 ISR 列表为空时(ISR 里的所有副本都挂了),是否允许一个不在 ISR 中的落后副本被选为新 Leader。
如果设成 true,会发生什么?
假设 Partition 0 有 3 个副本,Leader 在 Broker-0,两个 Follower 分别在 Broker-1 和 Broker-2。Broker-1 因为网络问题被踢出了 ISR,此时它的数据落后 Leader 500 条消息。
然后 Broker-0 和 Broker-2 同时宕机了。ISR 为空。
如果 unclean.leader.election.enable=true,Kafka 会让落后 500 条的 Broker-1 当新 Leader。那 500 条消息直接被回滚丢弃,消费者永远消费不到。
更隐蔽的是,这种丢失不会在生产者端体现出来,因为那 500 条消息当初发送的时候 ack 是成功的。只有事后对账才能发现数据少了。
# 确保这个参数是 false
unclean.leader.election.enable=false
设成 false 之后,ISR 为空时分区直接不可用,不会选出落后的 Leader。不可用虽然痛,但至少不会悄悄丢数据。挂了你知道,丢了你不知道,后者更可怕。
第四个坑:Page Cache 没刷盘,整机宕机数据蒸发
这个是最底层的坑,很多人根本没意识到。
Kafka 的写入路径是:消息先写进操作系统的 Page Cache,然后由 OS 在后台异步刷盘到磁盘。Follower 确认收到消息,也只是写进了自己的 Page Cache。
正常情况下这没问题,OS 会在几秒到几十秒内把 Page Cache 刷到磁盘。但如果 Broker 所在的物理机突然断电、内核 panic 或者硬件故障导致整机宕机,Page Cache 里还没刷盘的数据就彻底没了。
来看一下时序:
T0: Producer 发送消息 M1
T1: Leader(Broker-0) 写入 Page Cache,Follower(Broker-1, Broker-2) 也拉取并写入 Page Cache
T2: ISR 中 3 个副本都确认,ack 返回 success
T3: OS 还没来得及刷盘(Page Cache → 磁盘 的异步操作尚未执行)
T4: 机房 PDU 故障,三台 Broker 所在的机架同时断电
T5: 重启后,三台 Broker 的 Page Cache 全部丢失,消息 M1 不存在于任何副本的磁盘上
你可能觉得三台同时断电概率太低了。但同机架断电、同交换机故障这种事在大规模集群里不算罕见。而且不需要三台同时挂,只要 Leader 和 ISR 中的 Follower 恰好在同一个机架上,一个机架断电就够了。
Kafka 提供了两个参数来控制刷盘行为:
# 每收到 1 条消息就刷盘(极端保守,性能暴跌)
log.flush.interval.messages=1
# 每隔 1000 毫秒刷一次盘
log.flush.interval.ms=1000
但 Kafka 官方文档明确不推荐设这两个参数。原因很直接:强制刷盘会把 Kafka 的写入吞吐量从几十万 QPS 打到几千 QPS,性能下降一到两个数量级。
官方的建议是靠多副本跨机架部署来对冲 Page Cache 丢失的风险,而不是靠刷盘。
# broker 配置:指定机架信息
broker.rack=rack-1
配了 broker.rack 之后,Kafka 在分配副本时会尽量把同一个 Partition 的不同副本分散到不同机架上。这样即使整个机架断电,其他机架上的副本磁盘里大概率已经有数据了(因为刷盘的时机不完全同步)。
一个完整的安全配置长什么样
把前面说的坑都堵上,一个相对安全的 Kafka 消息可靠性配置是这样的:
# ---- Producer 端 ----
acks=all
retries=3
retry.backoff.ms=100
enable.idempotence=true
max.in.flight.requests.per.connection=5
# ---- Broker 端 ----
min.insync.replicas=2
unclean.leader.election.enable=false
default.replication.factor=3
broker.rack=rack-X
# ---- Topic 级别(可覆盖 Broker 默认值)----
# kafka-configs.sh --alter --topic order-topic \
# --add-config min.insync.replicas=2 \
# --bootstrap-server 192.168.1.10:9092
对应的 Java Producer 配置:
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "192.168.1.10:9092,192.168.1.11:9092,192.168.1.12:9092");
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.RETRIES_CONFIG, 3);
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5);
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
// 发送消息时一定要处理回调,不要 fire-and-forget
producer.send(new ProducerRecord<>("order-topic", orderId, orderJson), (metadata, exception) -> {
if (exception != null) {
// 发送失败,必须有兜底逻辑:写本地日志、入 DB 重试表、告警
log.error("消息发送失败, orderId={}", orderId, exception);
saveToRetryTable(orderId, orderJson);
}
});
这里有个细节:enable.idempotence=true 开启了生产者幂等。它能保证即使因为网络抖动触发了重试,同一条消息也不会在 Broker 端被重复写入。配合 acks=all 一起用,可靠性会高很多。
跨机架部署才是最后一道防线
前面说了,刷盘参数不推荐用,靠副本冗余来兜底。但副本冗余的前提是副本不在同一个故障域里。
如果你的 3 个副本恰好全在同一个机架上,那机架一断电,3 副本 = 0 副本,什么 acks=all、min.insync.replicas=2 全白搭。
配置 broker.rack 之后,可以用命令验证副本分布:
kafka-topics.sh --describe --topic order-topic \
--bootstrap-server 192.168.1.10:9092
# 期望输出:3 个副本分布在不同的 Broker 上,而这些 Broker 属于不同的 rack
# Partition: 0 Leader: 0 Replicas: 0,1,2 Isr: 0,1,2
# Broker-0: rack-1, Broker-1: rack-2, Broker-2: rack-3
如果发现副本集中在同一个机架,要用 kafka-reassign-partitions.sh 做一次副本迁移,把它们打散到不同机架上。
说到底,Kafka 的数据可靠性不是靠某一个参数保证的,是一套组合拳。acks=all 只是其中一拳,漏了任何一环都可能在极端场景下丢数据。
回到开头那个案例,排查到最后发现,那个 Topic 的 min.insync.replicas 是默认的 1,而且三个 Broker 恰好在同一个机架上。Broker-2 宕机的时候,另外两个 Broker 的 Page Cache 里有部分数据还没刷盘,运气差的那 1200 条消息就这么没了。
acks=all 保证的是"ISR 中的副本都收到了",但"收到了"不等于"落盘了","ISR 中的副本"也不等于"你以为的三个副本"。这两个不等号,就是消息丢失的全部真相。
我是小富,下期见。
