Kafka 的“消息不丢”不是一个配置能解决的,而是要在 生产者 → Broker → 消费者 三个环节同时做保证。
下面从“为什么会丢”到“怎么解决”做一次完整拆解。
先明确:Kafka 消息丢失发生在哪几层?
一条消息的生命周期:
Producer 发送 → Broker 存储/复制 → Consumer 拉取 → Consumer 提交 offset
所以丢失可能发生在:
生产者端:发送失败、确认策略太弱、缓冲丢失、重试不合理
Broker 端:副本数不足、ISR 门槛太低、unclean leader 选举、刷盘策略
消费者端:offset 提前提交、自动提交、rebalance 时处理不当
端到端语义:重复消费、消费逻辑非幂等、事务缺失
生产者
Kafka Producer(生产者)是向 Kafka Broker 发送消息的客户端,负责把业务数据封装成消息,发送到指定 Topic。核心流程:消息构建 → 序列化 → 分区选择 → 累加器缓存 → 发送线程批量发送 → Broker 应答确认。
核心架构组件
- 每个分区对应一个双端队列,队列元素是 ProducerBatch(批量批次)
- 一个 Batch 包含多条消息,达到大小或者时间阈值就触发发送
- 配置参数:
- buffer.memory:缓冲区总内存大小
- batch.size:单个批次最大字节数
- linger.ms:延时发送,没填满 batch 时,最多等待多久再发送,用来攒批量,提升吞吐量
分区规则:
Kafka 默认分区器:DefaultPartitioner
消费者必须使用对应的反序列化器。
发送消息三种方式
1. 异步发送(不关心结果,最常用)
producer.send(new ProducerRecord<>("topic", "key","value"));
直接返回,不等待 broker 响应,高吞吐;消息丢失风险取决于 acks 和重试配置。
2. 异步 + 回调(推荐,感知成功失败)
producer.send(record, new Callback() {
@Override
public void onCompletion(RecordMetadata metadata, Exception e) {
if(e != null){
// 发送失败
}else{
// 发送成功,拿到分区、offset
}
}
});
回调在 Sender 线程执行,回调逻辑不要写耗时操作,会阻塞发送线程。
3. 同步发送
调用 get() 阻塞等待结果:
RecordMetadata meta = producer.send(record).get();
吞吐量很低,适合强可靠场景。
核心可靠性参数
1. acks 应答级别
控制 broker 多少副本写入成功才返回成功给生产者。
- acks=0:生产者发完就认为成功,不等 broker 响应。吞吐量最高,可能丢消息。
- acks=1(默认):仅 leader 副本写入成功就应答。leader 刚写完挂掉,还没同步到从副本,消息丢失。
- acks=all /‑1:leader 和所有 ISR 内副本全部写入完成才应答。可靠性最高,吞吐量下降。
2. retries 重试次数
消息发送失败时生产者自动重试。配合 retry.backoff.ms 重试间隔。
注意:重试会带来消息重复,业务要做幂等处理。
3. max.in.flight.requests.per.connection
单个 broker 连接上,未完成的请求最大数量。
- Kafka <1.0:大于 1 时,重试会造成消息乱序;
- Kafka >=1.0 开启 enable.idempotence 幂等生产者,即使大于 1 也能保证顺序。
4. enable.idempotence 幂等生产者
默认开启 true。 解决:重试带来消息重复。原理:生产者分配 PID,每条消息带序列号,Broker 保存序列号,重复序列号直接丢弃。
幂等只能保证单会话、单分区不重复;生产者重启之后 PID 变化,无法跨会话幂等。
5. 事务生产者 transactional.id
用于跨分区、跨 Topic 的原子写入;要么全部成功,要么全部失败,配合消费者的 read_committed 隔离级别。
消息发送完整流程
- 如果当前 batch 未满,写入;
- 如果满了,新建 batch;
- 成功:释放 batch 内存,回调 onCompletion 成功;
- 失败:判断是否可以重试,放回累加器等待重试;不可重试则回调异常。
重要参数汇总
| bootstrap.servers | broker 地址列表 |
| key.serializer / value.serializer | 序列化器 |
| buffer.memory | 生产者缓冲区总内存 |
| batch.size | 单个批次大小 |
| linger.ms | 等待攒批时间 |
| acks | 副本应答策略 |
| retries | 重试次数 |
| enable.idempotence | 幂等生产者开关 |
| max.in.flight.requests.per.connection | 单连接并发请求数 |
| compression.type | 压缩算法:none、lz4、snappy、gzip,压缩提升吞吐量 |
常见问题
- acks=0/1,leader 宕机;缓冲区满阻塞 / 丢弃;生产者异常退出,缓冲区未发送消息丢失。 解决:acks=all,合理重试,设置回调捕获异常。
Kafka Producer(生产者)可能丢失消息的场景
前提:消息还没收到 Broker 成功 ACK,一旦生产者拿到成功 ACK,消息就已经持久化到 ISR 副本,生产者侧不再丢失。
场景 1:消息在 RecordAccumulator 缓冲区,还没发送,进程宕机
消息写入 ProducerBatch,还没满足发送条件(batch 未满、linger.ms 没到),仍然保存在生产者内存。 生产者进程直接崩溃、kill -9,内存数据全部清空,消息永久丢失。
- 常见:linger.ms 设置较大,攒批时间长;服务下线没有执行 flush () 和 close ()。
- 规避:优雅停机,调用producer.flush()再 close;业务层持久化消息,发送成功后再删除本地记录。
场景 2:acks=0,发送后不等 Broker 应答
生产者发出 TCP 数据包直接标记发送成功,完全不等待 Broker 返回结果。 网络丢包、Broker 宕机,Broker 根本没有收到消息,生产者却认为发送成功 → 消息丢失。
- 规避:高可靠场景禁止 acks=0。
场景 3:acks=1,Leader 写入成功、Follower 尚未同步,Leader 宕机
Leader 副本写入本地日志,立刻返回 ACK,此时 Follower 还没拉取同步这条消息。 Leader 机器故障宕机,集群从 ISR 选新 Leader,新 Leader 没有这条消息 → 消息丢失。
这个丢失发生在 Broker 侧,但根源是生产者应答级别配置过低。
- 规避:高可靠场景使用acks=all。
场景 4:缓冲区已满,send 阻塞超时抛出异常,消息丢弃
Broker 不可用,Sender 无法发送,累加器缓冲区持续堆积打满。 后续调用 send () 会阻塞,最长等待max.block.ms,超时抛出 TimeoutException,这条消息直接丢弃。 业务代码如果没有捕获异常、没有兜底保存消息,消息丢失。
- 规避:捕获超时异常,将消息存入本地 / 数据库死信,后续重试;合理调大 buffer.memory。
场景 5:不可重试异常,消息直接丢弃,不会重试
发生这类异常时,生产者不会重试,直接回调失败。如果业务不处理,消息丢失:
- 规避:Callback 捕获异常,不可重试消息写入死信队列。
场景 6:异步发送无 Callback,重试耗尽后静默丢失
业务代码直接producer.send(record)不传入回调。 多次重试之后仍然发送失败,Sender 线程抛出异常,业务代码完全感知不到,消息悄悄丢失。
这是生产中非常经典的坑。
- 规避:所有 send 必须携带 Callback,异常打日志、告警、落死信。
场景 7:flush () 之后立刻宕机
flush()的作用是把累加器里所有待发送 batch 交给 Sender 线程发起网络请求,不代表 Broker 写入成功。 flush 调用完成,但是网络请求还在路上、还没收到 ack,进程宕机,消息丢失。
很多人误解 flush = 保证落地,这是错误的。
场景 8:元数据拉取失败超时,消息无法写入
生产者没有 topic 元数据,请求 Broker 获取元数据长时间失败,send 阻塞 max.block.ms 超时,消息丢弃。
生产者消息重复、乱序场景
消息重复
Broker 已经成功写入消息,但是 ACK 包在网络传输中丢失;生产者超时触发重试,再次发送同一条消息 → Broker 收到两条一样的消息。
- 解决:开启enable.idempotence=true幂等生产者,PID + 序列号自动过滤重复消息;
- 限制:幂等生产者仅单会话、单分区生效,生产者重启后 PID 变化,无法跨会话防重复。跨分区原子写入需要事务。
消息乱序
Kafka 1.0 之前,没有幂等生产者,max.in.flight.requests.per.connection>1:
- 第 1 批消息发送失败,等待重试;
- 第 2 批消息发送成功;
- 第 1 批重试成功,消息顺序颠倒,产生乱序。
- 解决:开启幂等生产者;或者max.in.flight.requests.per.connection=1,同一连接只允许一个未完成请求。
生产者防丢总结
生产者消息丢失,基本都发生在消息拿到成功 ACK 之前。
高可靠生产者配置
acks=all
enable.idempotence=true
retries=2147483647
retry.backoff.ms=100
max.in.flight.requests.per.connection=5
buffer.memory=67108864
batch.size=16384
linger.ms=5
Broker端
Broker 端消息丢失,几乎全部集中在 "副本复制不完整" 或 "数据未落盘" 这两个根因上—— 只要 Leader 故障时,ISR 里没有一份完整副本可用,或所有副本的数据都还在内存页缓存里,消息就会丢。下面系统拆解。
说明:图为 Kafka 副本、应答门槛与丢失机制的通用结构示意,非特定集群测量结果。
Broker 端可靠性依赖的三根支柱(先理解,再谈丢)
一句话:只要 "选出的新 Leader 里没有完整数据",Broker 就丢消息。
Broker 端可能丢消息的完整场景
场景 A:acks=1,Leader 写完 ACK 后宕机(最经典)
- 根治:acks=all。
场景 B:acks=all,但 ISR 收缩到只剩 Leader 一个副本
- acks=all 的意思是 " 等 ISR 内所有副本写完 ",不是等全部副本。
- 根治:min.insync.replicas(topic 级)。设 2 表示 ISR 至少 2 份才允许写入;不足时直接拒绝写入(抛异常),宁可写不进去也不丢。
场景 C:unclean.leader.election.enable=true,选举非 ISR 旧副本当 Leader
- 参数默认 false(只能从 ISR 里选 Leader,ISR 全挂则该分区不可写,保可用性弃可用性)。
- 设为 true:允许选 数据不完整的非 ISR 副本 当 Leader。
- 根治:高可靠场景 关闭 unclean.leader.election.enable。
场景 D:ISR 内所有副本同时断电,数据仍在 PageCache 未落盘
- Kafka 写入:消息先写进操作系统 PageCache(页缓存) 就返回,不是立即刷盘;刷盘由 OS 后台或 log.flush.* 控制。
- 说明:只要有一个副本已刷盘且存活,选举它当 Leader 就不会丢。只有 ISR 全部机器同时断电才触发。
- 注意:不要靠频繁刷盘防丢(log.flush.interval.ms 调很小)会严重伤性能。Kafka 的可靠性靠副本,不靠刷盘。
场景 E:磁盘损坏 / 日志文件丢失(运维事故)
- 磁盘硬件故障、数据文件损坏、人为误删 topic 的 segment 日志。
- 若对应副本也被破坏,消息无法恢复。
场景 F:Follower 长期落后被踢出 ISR,副本数悄悄减少
- Follower 因网络抖动、GC 停顿长期追不上 Leader,被移出 ISR。
- 副本 "实际可用的" 数量下降。若此时 Leader 故障,可用的完整副本更少,丢失概率上升(本质是场景 B 的诱因)。
- 属于事前预警,不是立即丢,但会放大其它场景的风险。
Broker 端 "什么是安全的" 对照
| acks=0 | 最低 | 发出即成功,网络 / Broker 问题即丢 |
| acks=1(默认) | 中 | Leader 写入即成功,Leader 宕机会丢 |
| acks=all + 副本 3 + min.insync.replicas=2 + unclean=false | 高 | ISR≥2 才写,只从 ISR 选 Leader |
| 副本数 = 1 | 极低 | 无冗余,任何故障即丢 |
Broker 端防丢最佳实践(可直接背)
# Broker 端
default.replication.factor=3 # 默认副本数 3
min.insync.replicas=2 # topic 级:ISR 至少 2 份才写入
unclean.leader.election.enable=false # 禁止非 ISR 选主
# 生产者端配合
acks=all
enable.idempotence=true
- 生产环境禁止 auto.create.topics.enable=true 自动建单副本 topic(默认单副本,无冗余)。
- 多 Broker 跨机架部署,避免 ISR 副本集中在同一台物理机 / 同一机架(防场景 D)。
- 磁盘做好 RAID 冗余,降低场景 E。
高频面试题
Q:acks=all 就绝对不丢吗? 不是。① ISR 收缩到只剩 Leader,会退化成 acks=1;② ISR 全部同时断电且未刷盘;③ unclean 选举非 ISR 副本。所以 acks=all 必须配合 min.insync.replicas、unclean=false。
Q:为什么 Kafka 不靠频繁刷盘保证不丢? 因为 Kafka 的可靠性模型是 "多副本 + ISR",靠副本冗余而不是单机刷盘;强制频繁 fsync 会大幅降低吞吐,且仍无法抵御整机断电。只有 ISR 全机断电才丢,属于极端场景。
Q:unclean.leader.election.enable 怎么权衡? false = 保数据不丢,但 ISR 全挂时分区不可写(牺牲可用性);true = 保可用性,但可能丢数据。金融 / 订单类必须 false。
Kafka Consumer(消费者端)
核心一句话:消费者是主动拉取模型,消息丢失 / 重复几乎全部由 offset 提交时机决定,Broker 保存 offset 到__consumer_offsets系统主题。
一、消费者核心架构与基础概念
1. 消费模型
Kafka 消费者不是 Broker 推送消息,而是消费者主动调用 poll () 拉取。 消费者组 ConsumerGroup:同一个 group.id 下所有消费者构成一个消费者组。
- 一个分区,同一时刻只能分配给组内一个消费者(保证单分区顺序消费);
- 组内消费者数量 > 分区数:多余消费者空闲,没有分区;
- 组内消费者数量 < 分区数:一个消费者消费多个分区。
2. Offset(偏移量)
offset 是分区内消息的序号,代表消费位置。offset 持久化存储在 Broker 的__consumer_offsets系统 topic。
- 当前消费到哪条消息,由提交的 offset 决定;
- 提交 offset,代表 Kafka 认为这条消息以及之前所有消息已经消费完成。
重点:消息是否被处理成功,Kafka 本身完全不知道,只能相信消费者提交的 offset。 这是消费端所有问题根源。
3. 关键角色
二、消费者工作流程
三、Offset 提交的两种方式
1. 自动提交 enable.auto.commit=true(默认开启)
enable.auto.commit=true
auto.commit.interval.ms=5000 # 默认5秒自动提交
机制:poll 拉取消息后,每到时间间隔,自动提交上一轮 poll 拉取消息的 offset。 ✅ 优点:简单,不用手动写提交代码 ❌ 巨大风险,极易消息丢失 场景:poll 拿到消息,业务正在处理,还没处理完;5s 时间到,自动提交 offset。随后业务进程崩溃,这批消息已经提交 offset,Kafka 认为消费完成,重启消费者不会重新拉取 → 消息丢失
自动提交只适合:消息丢了影响很小的日志采集场景;订单、支付等强可靠业务严禁使用自动提交。
2. 手动提交 enable.auto.commit=false(高可靠业务推荐)
业务代码自己控制什么时候提交 offset,分两种:
最佳实践:消息全部处理成功之后,再提交 offset,保证 At-Least-Once(至少消费一次,可能重复,不会丢)。
四、消费者端消息丢失场景
消费者端天然不容易丢消息,错误的 offset 提交逻辑才会丢消息
场景 1:自动提交 offset,消息未处理完成,进程崩溃(最经典)
流程:
场景 2:手动提交,先提交 offset,后执行业务逻辑(代码写反)
// 错误写法!!
consumer.commitSync(); //先提交offset
processMessage(msg); //再处理业务
提交 offset 成功,业务执行失败 / 程序崩溃 → 消息丢失。
场景 3:重置 offset 到更早位置以外的错误位置
手动调用 seek 把 offset 设置到当前消息之后,跳过一批消息,直接丢失中间消息。
五、消费者端消息重复消费场景(高频考点)
At-Least-Once:Kafka 默认语义,消息至少被消费一次,大概率出现重复
场景 1:业务处理成功,但 offset 提交失败
业务逻辑执行完成,准备提交 offset,网络异常 / 进程宕机,offset 没有持久化到 Broker。消费者重启,从旧 offset 重新拉取,再次消费同一条消息。
场景 2:Rebalance 再均衡导致重复消费(最常见)
原因:rebalance 发生时,已经拉取、处理,但 offset 未提交。
场景 3:消费超时被踢出组(max.poll.interval.ms)
业务处理消息很慢,两次 poll 间隔超过max.poll.interval.ms,协调器判定消费者卡死,踢出组触发 rebalance,未提交 offset 消息重新分配,产生重复。
解决办法:业务处理放到独立线程池,poll 主线程只负责拉取消息、控制 poll 调用,不要阻塞 poll 线程;减小 max.poll.records,减少单次拉取消息数量。
六、消费者关键配置参数
| group.id | 消费者组标识,相同 group.id 属于同一组 |
| enable.auto.commit | 是否自动提交 offset,高可靠场景 false |
| auto.commit.interval.ms | 自动提交间隔 |
| max.poll.records | 一次 poll 最多拉多少条消息,消息处理慢要调小 |
| max.poll.interval.ms | 两次 poll 最大间隔,超时被踢出组,触发 rebalance |
| session.timeout.ms | 会话超时,心跳超时被踢出组 |
| heartbeat.interval.ms | 心跳发送间隔,一般是 session.timeout.ms 的 1/3 |
| auto.offset.reset | 分区无 offset 时的重置策略earliest:从头开始;latest:最新消息开始;none:抛异常 |
七、消费语义
Kafka 本身无法直接实现 ExactlyOnce,需要业务层配合。
八、消费者防丢 & 防重复最佳实践
九、面试简答
Q:消费者为什么会重复消费? A:业务处理成功,但 offset 未提交,发生 rebalance 或者消费者重启,新消费者从旧 offset 拉取消息,导致重复。
Q:消费端怎么实现不丢消息? A:关闭自动提交,消息处理成功之后手动提交 offset;保证业务处理成功,再提交 offset。
Q:auto.offset.reset 什么时候生效? A:分区没有保存该消费者组的 offset 的时候(第一次消费、offset 被删除),才会生效。已经有 offset,不会触发。