做Kafka运维和开发的这些年,我见过太多人一脸笃定地说“我配了acks=all,消息不可能丢”,结果线上数据还是对不上。也有不少团队在面试时把“Kafka为什么会丢消息”背得滚瓜烂熟,一遇到真实的丢数据告警就手忙脚乱。这个问题之所以经典,是因为它根本不在于某一个环节,而是一整条链路上每个看似合理的默认配置、每段没被认真处理的异常、每次想当然的“应该没问题”,都在悄悄给丢消息创造条件。
这篇文章我就基于自己的实战经验,把Kafka丢消息的根源按生产端、Broker端、消费端三个环节逐一拆开讲,再带大家走一遍我处理过的真实故障排查过程,最后给出一份可以直接对照检查的避坑配置。看完之后,你至少能回答三件事:自己的集群和客户端配置会不会丢消息、消息真的丢了该怎么查、以后怎么从设计上杜绝这个问题。
1. 生产端:消息在到达Broker之前就已经蒸发
很多人理解Kafka丢消息,第一反应是Broker挂掉导致数据丢失。但实际上,有很大一部分丢消息,消息压根就没离开过客户端所在的机器。生产端的丢消息往往最隐蔽,因为代码没有报错,日志没有异常,只是数据量对不上。
1.1 异常被吞掉:无辜的try-catch
我接手过一个业务团队的项目,他们在发送消息的代码里写了这样一段逻辑:
try { producer.send(new ProducerRecord<>("pay_order", orderId, message)); } catch (Exception e) { log.error("send message error", e); }表面上这没什么问题,异常打了日志,不会影响主流程。但Kafka生产者send()方法是异步的,它返回的是一个Future。你把消息丢进缓冲池就返回了,真正的网络发送是在后台线程完成的。如果只捕获send()方法同步抛出的异常,那些在后台线程发生的发送失败、超时、Broker不可达,都不会被你捕获到。
正确做法是给send()方法挂一个回调:
producer.send(new ProducerRecord<>("pay_order", orderId, message), (metadata, exception) -> { if (exception != null) { // 这里才是真正需要处理异常的地方 log.error("msg send failed, key={}, msg={}", orderId, message, exception); // 记录到本地文件、数据库,或发送到备用通道 } });我见过很多线上丢消息,最后查下来都是这种“send完就不管”的写法。业务方以为自己在发消息,实际上消息在客户端缓冲池里因为内存不足或者序列化失败被悄悄丢弃了。这里想提醒各位,凡是用了fire-and-forget(发后即忘)模式的发送,本质上就是主动放弃了Kafka的可靠性保障。
1.2 acks=all背后的坑:默认配置远没有想象中安全
生产端的核心参数是acks,它决定了生产者对消息写入成功的判定标准。
acks=0:不等待任何确认,消息发出即认为成功,这种配置下丢消息是必然的,它只适合日志采集这种允许丢的场景。acks=1:等待Leader写入本地日志即返回成功。这里有一个很多人没意识到的盲区——Leader写成功了,但是Follower还没同步,Leader一宕机,消息就没了。acks=all:等待ISR中所有副本都同步完成才返回成功。这才是我们通常认为的“不丢消息”配置。
但请注意,acks=all并不等于绝对安全。它有一个前提:你的min.insync.replicas配置必须大于1。我见过有人把acks配成all,但min.insync.replicas还是默认的1。这种情况下,如果ISR里只剩Leader一个副本,写入依然会被判定为成功。等Leader宕机,数据直接蒸发。这相当于你买了一把密码锁,却把密码设成了初始的0000。
acks=all会不会影响性能?会,但现在的Kafka生产者默认开启了幂等(enable.idempotence=true),配合高版本协议,acks=all的开销已经比早期版本小很多。在追求数据可靠性的场景里,这个性能代价完全值得。
1.3 重试与超时:你以为在重试,实际上已经放弃
生产端的另一个隐蔽问题是重试参数配置不当。Kafka客户端默认的retries是Integer.MAX_VALUE,从参数上看是无限重试,但实际重试行为受delivery.timeout.ms控制,这个参数默认是120秒。也就是说,一条消息在120秒内没发送成功,就直接判失败,重试再多也没用。
我曾经排查过一个案例:某个业务高峰时段,下游消费变慢,导致Broker端写入耗时上升,生产者客户端一批批消息发送超时。他们的retries设了个很大的值,以为会一直重发,但delivery.timeout.ms没调大,所以大量消息在重试几次后直接超时失败。更麻烦的是,他们的代码用的是前面说的“发后即忘”模式,根本不知道有消息失败。
这里给一个实操建议:生产环境一定要把delivery.timeout.ms适当调大,比如设为5分钟以上,同时配合max.block.ms,给消息发送留足缓冲时间。另外,重试本身会引入消息乱序的问题。如果你依赖Kafka分区保证消息顺序,并且开启了重试,建议设置max.in.flight.requests.per.connection=1(如果开启了幂等,可以不用调成1,Kafka会自动处理),避免因为前一条消息重试、后一条消息已经发送导致顺序颠倒。
1.4 大消息被静默拒绝:1MB限制与InvalidReceiveException
有一个热搜词是“kafka 接收1m”,对应的就是一个高频问题:消息体超过1MB怎么办。Kafka Broker端默认的message.max.bytes是1048576字节,也就是1MB。如果你发的消息超过了这个值,客户端会收到RecordTooLargeException,这条消息直接被判定为发送失败。
更麻烦的是Broker端的InvalidReceiveException。这个异常我在生产环境遇到过好几次,它出现在客户端发送超大请求、Broker端解析请求头时发现消息大小和它预期的缓冲区不匹配。如果你把客户端的max.request.size调大了,比如调到10MB,但Broker端的socket.request.max.bytes没跟着调大,客户端和服务端对最大请求体大小的认知就不一致。Broker解析不了这么大的请求,会直接抛出InvalidReceiveException并断开连接,消息还没进入Kafka就没了。
经验之谈:传输超过1MB的消息,不建议一味调大Kafka的限制参数。消息太大意味着在网络上传输时间更长、Broker端落盘和复制的时间更长、消费者拉取的带宽占用也更大,整个链路的故障概率都会上升。更合理的方案是:把大消息的Payload传到对象存储或HDFS,Kafka里只保存这个消息的索引信息。如果实在要传,就同步调整这几处参数:
| 参数 | 默认值 | 作用位置 | 建议值(传5MB消息时) |
|---|---|---|---|
max.request.size | 1MB | 生产者 | 6MB |
message.max.bytes | 1MB | Broker | 6MB |
replica.fetch.max.bytes | 1MB | Broker(副本同步) | 6MB |
fetch.max.partition.fetch.bytes | 1MB | 消费者 | 6MB |
这几处必须联动调,少配任何一处,都会在某一环把大消息卡死。
2. Broker端:写入“成功”并不等于数据“安全”
如果你确认生产端代码没问题、消息确实发到了Broker,接下来就要审视服务端本身。Broker端丢消息的场景主要集中在副本机制配置、Leader选举策略和日志管理三个方面。
2.1 单副本的诱惑:Kafka默认就是单副本
Kafka创建Topic时,默认的replication.factor是1。这意味着你的Topic只有一个副本,这个副本本身就在Leader节点上。只要这个节点宕机,或者磁盘坏了,这个Topic的所有数据就彻底没了——包括那些已经返回“写入成功”的消息。
很多第一次搭建Kafka集群的人会踩这个坑:用默认配置创建Topic,数据量小的时候一切正常,等节点宕机或者做滚动重启时,发现部分Topic的数据不翼而飞。我强烈建议在服务端配置里加上default.replication.factor=3。同时配合Broker级别的min.insync.replicas=2,这样即使某个分区的Leader节点意外宕机,其他节点上还有完整的副本数据,可以正常接管。
2.2 ISR收缩:Leader明明活着,消息却已经不安全
ISR(In-Sync Replicas)是Kafka副本机制的核心概念。它指的是和Leader保持同步的副本集合。如果ISR里有Leader和一个Follower,写入就需要同步到Follower才算成功,这就是acks=all的含义。
但ISR是会收缩的。如果某个Follower同步速度太慢,或者网络波动导致它和Leader的连接超时,它会被踢出ISR。如果你的min.insync.replicas=1,当ISR收缩到只有Leader自己时,写入依然成功。但这时你已经没有任何实时备份了,数据处于裸奔状态。
这里有一个反直觉的细节:acks=all在ISR正常时有很好的保障,但在ISR收缩到临界值时,它反而是最危险的时候。因为客户端拿不到“写入失败”的反馈,会持续往一个没有备份的Leader上写数据。一旦Leader出问题,这段时间写入的所有数据全部丢失。所以一定要设置min.insync.replicas=2,让客户端在ISR不足时直接发送失败,你宁可让业务报错,也不能让数据不声不响地丢。
2.3 Leader选举的脏数据:unclean选举是丢消息的大户
这一节要解决一个非常经典的认知误区:Kafka丢消息,很多情况下不是因为数据没写入,而是因为Leader切换时选出了一个“落后的副本”当新Leader,然后旧Leader上的数据被直接截断。
Kafka里有一个参数叫unclean.leader.election.enable,默认是true。这个参数决定了当ISR中所有副本都宕机时,Broker是否允许一个不在ISR中的副本(即和Leader数据差距很大的“脏副本”)被选举为新Leader。
- 允许(
true):优先保证可用性,但脏副本成为Leader后,它缺失的数据会被视为“从未存在”,其他副本还会以它为基准进行日志截断。数据彻底丢失。 - 禁止(
false):优先保证一致性,宁可让分区不可用,也不丢失数据。等原来的Leader恢复,它带着完整数据回来,再继续对外服务。
我在生产环境见过一次严重的丢数据事故,就是unclean.leader.election.enable=true导致的。当时某个分区所在的机器磁盘故障,ISR里的副本都跟着出了异常,Broker自动把一台落后的副本拉起来当Leader。结果那个Topic当天的数据少了一大截,业务方还以为是上游没发消息。
另外还要提一下leader.epoch机制。Kafka从0.11版本引入了leader epoch,用来防止“旧Leader复活后,因为对已经提交数据的错误截断而导致的数据丢失”。但是注意,leader epoch解决的是新旧Leader之间的数据截断问题,它不解决unclean选举带来的“脏副本上位”问题。这两件事要分清。
2.4 磁盘物理损坏:日志段删了就真的没了
最后补一个很多人不重视的物理层面问题。Kafka的数据是落到磁盘上的,落盘机制依赖操作系统的PageCache来提升性能。数据先写PageCache,由操作系统异步刷到物理磁盘。如果机器突然断电,PageCache里还没来得及刷盘的数据就会丢。
Kafka提供了log.flush.interval.messages和log.flush.interval.ms两个参数来控制刷盘频率。但我要提醒的是:不要为了“防止丢数据”而把刷盘频率调得特别高。这会让Kafka的吞吐量直线下降,而且这两个参数在生产环境里一般不建议动,让操作系统自己决定刷盘时机就够用。真正需要做的是基础设施层面的保障:使用RAID阵列、配置UPS电源、定期做备份。Kafka不是数据库,它不承担数据最终兜底的职责,你的安全底线应该建立在多副本和备份策略上。
3. 消费端:最难察觉的丢消息其实在业务代码里
如果说生产端和Broker端的丢消息还能通过监控和日志去定位,那消费端的丢消息就完全是“剪不断理还乱”。因为消费端丢消息的表现不是“消息没了”,而是“消息还在,但业务状态已经错了”。
3.1 自动提交offset:消息还没处理完,位移已经报账了
enable.auto.commit的默认值是true。这意味着消费者客户端会每隔auto.commit.interval.ms(默认5秒)自动提交一次消费位移。很多新手写消费逻辑是这样的:
while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecord<String, String> record : records) { saveToDatabase(record); // 正常业务处理 } // 这里什么都没做,offset由Kafka自动提交 }看起来没毛病,处理完了就提交。但自动提交的时机是客户端确定的,它不关心你是刚拿到消息还是已经处理完。举个例子:poll()拉回来100条消息,你处理到第50条时进程崩溃,那么正常的预期是消费位点应该停在还没处理的第50条位置。但由于自动提交已经在某个时间点把位点提交到了第100条,进程重启后,从第51条到第100条的消息就永远不会被消费了。
这就是“数据没处理好,但offset却认为已经处理完”——消息直接消失。处理办法很简单:把enable.auto.commit设为false,改为手动提交。
3.2 手动提交的学问:先处理业务再提交位移
手动提交也不是随便提交就行的。核心原则就一条:必须先保证业务处理成功,再提交位移。
建议的流程是这样:
while (true) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecord<String, String> record : records) { process(record); // 业务处理 } consumer.commitSync(); // 全部处理成功后再提交 }这个例子里的关键点在于:commitSync()确认的是前面所有消息都已经成功处理,然后才移动消费位点。如果process()中间抛了异常,位点不会提交,消息会再次被消费。很多人觉得这样会导致重复消费,对,确实会重复。但你要想清楚一个优先级问题:重复消费是可以被业务幂等设计消化的,消息丢失却是无法找回的。两害相权,宁可重复,不可丢失。
我看到实际生产里还有一种折中写法:为了保证性能,一部分人用commitAsync()异步提交。这个API和commitSync()的区别在于,异步提交不会阻塞主线程,但它并不保证提交一定成功。如果提交失败,就会导致重复消费。我的建议是:commitAsync()可以配合commitSync()一起用,正常情况用commitAsync()提高吞吐量,应用关闭前最后再用一次commitSync()兜底。
3.3 多线程消费与消息顺序性:热搜里那个“技术债”
热搜词里有一条“kafka消费端多线程如何保证消息顺序性”,这道题我在面试里也经常问。先理清Kafka顺序性的底层逻辑:Kafka只保证单分区内消息有序,消费端如果用单线程处理消息,天然有序。一旦你把多线程引进来,顺序性就会因为并发处理而被打乱——订单创建消息在订单支付消息之后被处理,这种事故是会真实发生的。
常见的多线程消费设计有三种:
- 单线程
poll(),多线程process():这是最简单的并发方案,吞吐量提升明显,但消息顺序完全无法保证。适合对顺序不敏感的业务,比如日志处理、通知推送。 - 单线程
poll(),按分区key把消息分发到对应的工作线程队列:比如有10个分区,你用10个线程,每个线程维护一个队列,poll()出来以后按消息所属分区放进对应队列。同一个分区的消息始终由同一个线程处理,顺序性得到保证,吞吐量也能提升约等于分区数量的倍数。 - 一个分区一个消费线程:每个线程独立
poll()。但这个方案会受限于分区数量,分区就那么多,线程数很难弹性扩展。
如果你的业务强依赖顺序,像订单状态流转、库存扣减,我建议用方案二。注意这里有个容易踩的坑:方案二要求你的消息分发逻辑必须严格按partition维度来做,而不是按key做哈希。如果同一条业务线的消息被哈希到不同线程的队列,顺序还是乱。poll()返回的ConsumerRecord本身就带partition()字段,直接用它做分发最可靠。
再说回丢消息这件事。多线程消费之所以和丢消息有关系,是因为多线程处理时,“处理成功”的判定时机更容易出错。比如线程处理完业务但还没来得及提交offset,另一个线程触发了Rebalance,已经处理完的消息被迫再次消费,可能引发重复写入。虽然这不是“丢”,但它和丢消息同根源:offset提交时机和处理结果之间没有保证一致。
3.4 Rebalance风暴:消费组在壮烈中弄丢位移
Rebalance是Kafka消费者组的自愈机制,但它也是一把双刃剑。当消费者实例发生变动(新增、宕机、心跳超时)时,Kafka会触发Rebalance,把分区重新分配给存活的消费者。这个过程中,如果你的位移提交异常,就可能出现消费位点回退或前进的情况。
最典型的一个丢消息场景是这样的:消费者A在处理分区P的消息时,因为处理太慢导致心跳超时,被消费组判定为宕机,触发Rebalance,分区P被分配给消费者B。消费者A在被踢出前已经处理了一些消息,但还没来得及提交位移。注意,这里不是直接丢,而是消费者B会从旧的位移开始重新消费——重复。
反过来还有一种更隐蔽的丢失:消费者A的位移已经提交了,但它其实还有个本地缓冲区里没处理完的消息。这时Rebalance触发,它带着没处理完的数据“死了”,这些数据就彻底丢了。在高并发消费场景里,poll()拉取的消息往往比实际处理的速度快,消息会暂存在内存里。这批“在途消息”一旦遇到进程崩溃或者Rebalance,就凭空消失。
要降低这个风险:一是把max.poll.records调小一些,比如500改成100,减少单次poll()拉取的在途消息量;二是调大max.poll.interval.ms,给业务处理留足时间,避免因为慢消费被误判为宕机;三是最关键的,永远不要在消费线程里做耗时操作,比如同步调外部接口,这几乎必然导致Rebalance频繁触发。
4. 一场真实故障:从收到告警到定位根因的全过程
理论说了这么多,可能还是不如一个完整的排查链路有说服力。我挑一个印象最深的线上事故来讲,当时从头到尾用了两个多小时才把根因挖出来,整个过程基本涵盖了前面提到的所有环节。
4.1 故障表象:数据报表对不上账
那是一个交易系统的数据同步链路,上游服务通过Kafka发送订单消息,下游是报表统计服务。某天下午,业务方突然反馈:当天的订单量和报表统计数差了将近8000条。第一反应是下游统计逻辑有bug,但核对代码后没发现问题。于是开始查链路。
4.2 逐层排查:生产端、Broker端、消费端三方印证
第一步,先看监控。Kafka的JMX指标里有一个record-error-rate(发送失败率),拉出来看当天的峰值,发现确实有段时间这个指标不为零,集中在下午2点到2点20分之间,数值不算高但真实存在。同时生产端的日志文件里,用自定义回调打出的错误日志也有记录,说明在消息生产端就已经有异常发生了。
第二步,看Broker端的日志。在错误日志里发现了一个关键词:NotLeaderOrFollowerException。这个异常表示消息发送时,客户端找不到分区的Leader副本。结合时间点,那段时间正好Kafka集群有一台节点在做滚动重启。这个发现解释了一部分发送失败的原因,但依然不能解释全部——因为生产端有重试机制,短暂的Leader切换一般会重试成功。
第三步,看消费端。这是整个事故最关键的转折点。我们拉出了那个Topic的消费组Lag监控,发现消费组在下午2点左右出现了Lag断崖式下跌,从持续稳定的某个值突然跌到接近零。这说明消费位点在那个时间点被强制推进了。
再核对消费端的启动参数,发现enable.auto.commit=true,auto.commit.interval.ms=5000。也就是说,消费位点每5秒自动提交一次。那段时间因为Broker重启,生产端发送失败,部分消息根本没有成功进入Kafka,消费端自然也就拉不到这些消息。但由于自动提交的存在,消费位点并没有因为消息缺失而停在原地,而是顺着现有的消息继续往前了。
这下整个逻辑就完全闭合了:生产端发送失败(因为Leader切换)导致部分消息没发出去,中间又没有足够好的重试补偿;消费端自动提交offset,对缺失的消息毫无感知,也没有发生长时间Lag,所以监控没有报警。两个环节各自的小问题叠加,造成了一次完整的数据丢失。
4.3 修复方案:三管齐下
事故处理完,我们做了三处修改:
生产端:把
acks从默认值改为all(其实是-1),把retries保持一个较大的值,同时把delivery.timeout.ms从默认的120秒调到300秒,并且在send()方法的回调里增加失败消息本地落盘,写好后会有一个补偿Job定期从这个本地文件里重新发送。这一步保证了消息但凡进了生产者缓冲池,就尽可能被送达。Broker端:确认所有Topic的
min.insync.replicas为2,default.replication.factor=3,unclean.leader.election.enable=false。消费端:
enable.auto.commit改为false,改成业务处理成功后调用commitSync(),并且给所有下游消费逻辑加上了幂等处理。即使极端情况出现重复消费,也只是多执行一遍,不会产生脏数据。
这三个修复做完后,类似问题再没有发生过。回想整个排查过程,最耗时的地方其实不是定位参数配置,而是确认“消息究竟是在哪一个环节丢的”。生产端、Broker端、消费端三方的参数互相影响,如果不按链路逐步排查,很容易被表象误导。
5. 避坑配置手册与Kafka的真实边界
最后整理一份可以直接拿去对照检查的配置清单。我给每个配置项都标注了它对应的丢消息风险,方便业务侧做自查。
5.1 快速对照检查:你的配置安全吗?
| 配置项 | 推荐值 | 默认值 | 不配置的风险 |
|---|---|---|---|
acks | all | 1 | Leader宕机导致已写入消息丢失 |
min.insync.replicas | 2 | 1 | ISR只剩Leader时写入成功但无备份 |
replication.factor | 3 | 1 | 单副本无容灾能力 |
unclean.leader.election.enable | false | true | 脏副本当选Leader导致旧数据被截断 |
enable.auto.commit | false | true | 消息未处理完但offset已提交 |
delivery.timeout.ms | 300秒以上 | 120秒 | 消息在客户端重试超时被放弃 |
max.poll.records | 100~500 | 500 | 在途消息过多,Rebalance时丢内存数据 |
max.poll.interval.ms | 根据处理耗时调大 | 300000 | 慢消费导致被踢出消费组 |
5.2 消息延迟高不等于丢消息,但要警惕联动风险
热搜词里有一条“Kafka消息延迟高”,这里顺便说一句。很多人把消息延迟和消息丢失混为一谈。延迟高的原因很多:生产者linger.ms设置得大、消费者处理速度慢、分区数不足、Broker磁盘IO饱和等。这些本身不直接导致丢消息,但延迟高会让消费端的处理积压,进而触发max.poll.interval.ms超时、Rebalance、Commit失败等连锁反应,最终演变成丢消息或者大规模重复消费。
我见过一个团队,因为下游消费慢导致Lag持续上涨,他们把max.poll.interval.ms调到了10分钟,以为能缓解问题,结果消息积压超过这个阈值后,消费者还是被判定为失联,触发Rebalance,反而加剧了问题的复杂性。处理消息延迟的正确姿势是找出瓶颈:是下游数据库慢?是处理逻辑有串行依赖?还是分区数不够导致消费并发度上不去?单纯调Kafka参数往往是吃力不讨好。
5.3 检测丢消息的可观测手段
丢消息这种事,最怕的就是发生后很长时间没人发现。对于核心链路,我建议做三件事:
- 记录生产端的成功和失败指标,重点看
record-error-rate和回调里的异常日志。 - 对消费组持续监控Lag,Lag出现异常的断崖式下跌,往往比上涨更难排查,值得设置独立告警。
- 在业务侧做发送与消费数量的对账。最简单的实现是:生产端把消息ID写入Redis的Set中(带TTL),消费端消费成功后也写入同一个Set,然后用定时任务比对两边数量,超过阈值就告警。在高并发场景下这个方案代价不小,需要评估和取舍,但对核心账户、订单类数据链路,这个成本是值得花的。
5.4 理解Kafka的能力边界
最后我想聊一个理念问题。很多初学Kafka的人对它有一个误解:认为Kafka是配置好就绝对不会丢消息的。实际不是这样。Kafka的定位是一个高性能的分布式消息流平台,它的可靠性是建立在“合理配置+使用者正确操作”基础上的。它能做到不丢消息的边界是:生产端正确配置acks=all并处理好发送失败;Broker端副本数足够且ISR机制正常工作;消费端正确处理offset提交并设计好异常恢复流程。这三者任何一个环节被破坏,Kafka都会毫不犹豫地丢消息。
反过来看,Kafka有一些丢消息其实是“设计使然”。比如log.retention.hours默认168小时,消息超过7天就会被删除。比如启用日志压缩的Topic,key相同的新消息会覆盖旧消息。这些不是故障,是Kafka的功能特性。你要做的不是在报警时纠结“为什么消息没了”,而是在设计阶段就想清楚每条数据的生命周期和可靠性等级,然后对应配置。该用Kafka的地方用Kafka,该用数据库的地方你别偷懒,该加备份的加备份,该做对账的做对账。
每次踩过Kafka丢消息的坑之后我都会跟团队说一句话:Kafka本身不丢消息,丢消息的往往是使用它的人忽略了某个默认值。把默认值当成设计意图,是分布式系统使用里最贵的一堂课。希望看到这篇内容的朋友,都能带着这份对照清单,回去好好检查一遍自己的Producer和Consumer配置,别等数据对不上的那天再来翻文档。