Skip to content

Kafka 消息丢失场景与排查:生产端/消费端/服务端全链路分析

提出问题

"Kafka 到底会不会丢消息?"——这是面试官最爱问的送命题,也是生产环境最常踩的坑。Kafka 宣传的"高性能"恰恰是消息丢失的根源:为了吞吐它默认不刷盘、不等待确认就用异步方式返回。更棘手的是,消息丢失可能发生在生产端(发送时没确认)、Broker 端(副本还没同步就宕机)、消费端(offset 提交了但业务没处理完)三个环节,每段链路都有一堆配置参数暗藏玄机。

面试官问这个问题,不是让你背配置——他在确认你有没有亲手在线上踩过坑,并且知道怎么修。下面从三段链路逐一拆解丢失场景和排查手段。

分析问题

生产端丢失:配置没整对,发了等于没发

生产端消息丢失是最容易被忽视的环节。原因很直接:Producer 默认 acks=1,只要 Leader 写入本地日志就返回成功,但此时 Follower 还没同步,Leader 一旦宕机这条消息就丢了。

java
// 生产端安全配置:acks=all + 重试 + 幂等
Properties props = new Properties();
props.put("bootstrap.servers", "broker1:9092,broker2:9092,broker3:9092");
props.put("acks", "all");                          // 等待所有 ISR 副本确认
props.put("retries", Integer.MAX_VALUE);           // 无限重试
props.put("max.in.flight.requests.per.connection", "5"); // 幂等下可 >1
props.put("enable.idempotence", "true");           // 幂等 Producer(防重复)
props.put("delivery.timeout.ms", "120000");        // 整个发送超时 2 分钟
props.put("request.timeout.ms", "30000");          // 单次请求超时 30 秒

KafkaProducer<String, String> producer = new KafkaProducer<>(props);

丢失场景列举:

  • acks=0:Producer 发完就丢,不管死活。日志场景可能用,业务场景绝对禁止。
  • acks=1 + 无重试:Leader 写入后、Follower 同步前宕机,消息丢失且 Producer 收到成功回调,业务方毫不知情。
  • 重试次数不足retries=0retries=3 遇到 Broker 端短暂抖动(如 GC 暂停),重试耗尽后返回异常,但业务代码可能直接吞掉异常。

排查工具:kafka-producer-perf-test.sh 可以模拟生产压测看错误率;Producer 端抓日志搜 WARNERROR 级别的 org.apache.kafka.clients.producer

Broker 端丢失:Page Cache 是一把双刃剑

Broker 端丢失是 Kafka 最核心的"设计缺陷"——Kafka 不主动刷盘,数据写入 Page Cache 就返回成功,真正落盘依赖操作系统后台回写(pdflush)。断电或进程崩溃时,未刷盘的数据全部丢失。

bash
# 查看当前系统的刷盘参数
$ cat /proc/sys/vm/dirty_ratio
20
$ cat /proc/sys/vm/dirty_background_ratio
10
$ cat /proc/sys/vm/dirty_expire_centisecs
3000

dirty_ratio=20 意味着 Page Cache 中的脏页达到总内存 20% 时才触发同步回写,期间宕机数据就丢了。

Broker 端其他丢失场景:

  • ISR 收缩:某个 Follower 因 GC 暂停或网络抖动落后,被踢出 ISR。此时如果 min.insync.replicas=1(默认),acks=all 实际退化为 acks=1,Leader 宕机就丢数据。
  • Unclean Leader Electionunclean.leader.election.enable=true 时,ISR 全挂后允许 OSR 副本成为 Leader,该副本落后 Leader 的数据被"截断",消息丢失。
  • 日志过期删除log.retention.hourslog.retention.bytes 触发日志段删除,如果 Consumer 消费速度跟不上,未消费的消息就没了。

排查方法:

bash
# 检查 ISR 是否完整(Under-Replicated 分区数 > 0 表示有副本落后)
$ kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic my-topic
Topic: my-topic    Partition: 0    Leader: 1    Replicas: 1,2,3    Isr: 1,2
# 注意:ISR 只有 1,2,缺少 3,说明副本 3 落后了

# 查看 Broker 的磁盘使用率
$ df -h /data/kafka
Filesystem      Size  Used Avail Use% Mounted on
/dev/sda1       500G  480G   20G  96%  # 磁盘快满了,可能导致副本被踢出 ISR

消费端丢失:offset 提交时机是关键

消费端丢失是最常见的线上事故原因——原因不是 Kafka 的错,是消费代码写得有问题。

java
// 错误示范:先提交 offset,再处理业务
consumer.subscribe(Arrays.asList("my-topic"));
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    consumer.commitSync();  // ⚠️ 先提交 offset
    for (ConsumerRecord<String, String> record : records) {
        processRecord(record);  // 如果处理时宕机,这条消息丢了
    }
}
java
// 正确写法:处理完业务再提交
consumer.subscribe(Arrays.asList("my-topic"));
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    for (ConsumerRecord<String, String> record : records) {
        processRecord(record);
    }
    consumer.commitSync();  // ✅ 处理完再提交
}

消费端丢失的典型场景:

  • 自动提交 (enable.auto.commit=true):默认每 5 秒自动提交一次,两次自动提交之间宕机,已经 poll 但还没处理的消息再也拿不到了。
  • 手动提交但提前提交commitSync() 调用在 processRecord() 之前,宕机时消息丢失。
  • 异步提交 + 忽略回调commitAsync() 不传回调,提交失败不重试,offset 回滚导致下次消费重复(不是丢失,是重复;但配合错误处理可能导致"业务"层面的丢失)。

排查工具:

bash
# 查看 Consumer Group 的 offset 和 Lag
$ kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group my-group --describe
GROUP           TOPIC           PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
my-group        my-topic        0          1500            2000            500
# LAG=500 表示有 500 条消息未消费
# 如果 CURRENT-OFFSET 增长很快但 LAG 不降,说明 Consumer 可能跳过了消费

总结

消息丢失的全链路解决思路可以归纳为三段防御

环节关键配置/做法常见失误
生产端acks=all + retries=MAX + enable.idempotence=trueacks=1 以为安全了;重试次数设太小
Broker 端min.insync.replicas=2 + replication.factor=3 + unclean.leader.election=false默认 min.insync.replicas=1 形同虚设;Page Cache 不刷盘
消费端enable.auto.commit=false,处理完业务手动 commitSync()先提交再处理;commitAsync() 回调不加重试

面试话术示例:"Kafka 消息丢失分三段排查。生产端我遇到过 acks=1 导致 Leader 宕机丢消息,改为 acks=all + min.insync.replicas=2 解决;Broker 端最坑的是 Page Cache 不刷盘,我们加了 flush.messages=1 做双保险;消费端同事踩过先提交 offset 再处理业务导致丢消息的坑,通过代码 review 发现并修复。最终在业务层还加了对账机制——Producer 发消息前记录到数据库,Consumer 处理后发送回执,定时对账补发,做到业务层面的不丢不重。"

参考:Apache Kafka 官方文档 (https://kafka.apache.org/documentation/);《Kafka 权威指南(第 2 版)》Neha Narkhede 等

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