拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Kafka丢消息排查指南:从acks、ISR到偏移量的可靠性配置全解析

Kafka丢消息排查指南:从acks、ISR到偏移量的可靠性配置全解析

Kafka 在很长一段时间里都被当成消息中间件里的“天花板”来用,吞吐量高、可扩展性好、还能做存储。但凡是线上实际用过 Kafka 的团队,几乎都经历过一次“到底丢没丢消息”的争论。有人说是 Kafka 的问题,有人说是自己代码的问题,还有人说是网络抖动导致的偶发故障。实际上,我做了这么多年 Kafka 相关的项目,结论非常明确:Kafka 本身很少会“无理由丢消息”,绝大多数丢消息的根源,都出在生产者配置、消费者提交位点的时机,以及你对确认机制的理解上。

这篇文章我会从生产端、Broker 端、消费端三个链路分别拆解,把“丢消息”这件事彻底讲透。不管你是在做日志采集、实时数仓,还是普通的后端消息推送,这篇内容都值得对照你当前的配置逐条查一遍。你会发现很多丢消息不是玄学,而是配置和代码写得不够谨慎。

1. 先说结论:丢消息从来都不是 Kafka 一个环节的问题

1.1 我对“Kafka 丢消息”这句话的理解

很多人一搜“Kafka 丢消息”,第一反应是 Kafka 这个中间件本身有没有 bug。其实 Kafka 作为分布式系统,单从设计上讲,它已经尽最大努力去保证消息不丢了。问题在于,这个“保证”是有前提条件的,比如生产者要设置合理的确认机制、Broker 端要有足够的副本数、消费者要自己管理好偏移量。

你如果把保证可靠性的责任全部甩给 Kafka,自己代码里毫不在意地往生产端一塞、消费端一拉,那丢消息几乎是必然的。所以我的第一句话永远是:Kafka 不丢消息是有条件的,你要做的不是问 Kafka 为什么丢,而是问自己在哪个环节破坏了它的可靠性边界。

1.2 先想清楚:消息从进到出,要经过哪几道门

一条消息从业务系统产生,到最终被下游消费者处理,完整的路径是:

  • 生产者把消息发送给 Kafka Broker;
  • Broker 把消息写入分区对应的日志文件,并完成副本同步;
  • 消费者从 Broker 拉取消息,处理业务逻辑;
  • 消费者提交偏移量(offset),确认这条消息已经被消费。

这里面每一道门都有可能把消息弄丢。生产端可能因为发送失败、回调没处理、缓冲区满了而丢;Broker 端可能因为刷盘不及时、副本同步失败、Leader 切换而丢;消费端可能因为位点提交太快、程序崩溃、重平衡而丢。

你光盯着某一端排查是没用的。我见过不少团队 debug 了好几天,最后发现是消费端业务异常导致位点已经提交了但消息没处理完成,这种“假成功”最容易让人误以为是 Kafka 在丢消息。所以先建立这个全局链路认知,后面排查起来才有方向。

2. 生产端:你还没发出去,消息就已经丢了

2.1 acks=0 的代价:说发就发,后果自负

生产者端第一个经典坑就是 acks 配置。如果你把acks设成 0,意味着生产者不需要等待 Broker 的任何确认,消息发出去就算完事。这个配置下吞吐量最高,但代价就是,只要网络发生抖动、请求还没到 Broker,或者 Broker 写入过程中出现问题,生产者这边完全无感知。

我实际帮团队排查过一个采集系统的丢数据问题,客户端采集到日志后直接往 Kafka 里怼,配置里acks=0,结果高峰期数据丢失率在 0.5% 左右,看着不高,但对于需要精准对账的业务来说完全不可接受。后来把acks改成 1,丢失率立刻降下来了。

在这个配置下,你所谓“发消息成功”只是客户端把数据交给了 socket 缓冲区,并不是真的写到了 Kafka。所以生产环境里除了极少数允许丢失的监控指标类数据,我认为一般不要用acks=0。真到了丢不起的场合,再回来的查证成本远高于那点吞吐量收益。

2.2 acks=1 的中间地带:Leader 挂了你就傻眼

acks=1是很多团队里的默认选择,意思是生产者只需要等 Leader 副本写入成功就返回成功。这个配置在实际吞吐量和可靠性之间确实比较均衡,很多业务也都跑在这种模式下。但你要清楚它的死角:如果 Leader 在写入完成后、还没等其他副本同步的时候挂了,这条消息对消费者来说就是“丢失”的。

举个例子,你给主题配置了一个分区三个副本,某台机器上的 Leader 副本把消息写盘后,数据还没同步到其他两个 Follower,结果这台机器突然断电。Kafka 会从剩下的 Follower 里选一个新的 Leader,之前那批没有同步完成的消息就永远没了。

在这个环节,你要是只盯着生产者有没有收到成功回调,那是不够的。acks=1只保证“那一刻的 Leader 写入成功”,不保证这条消息在集群里永久存活。对于账务类、订单类这种强一致场景,我建议直接上acks=all,别太抠那一点点延迟。

2.3 最稳妥的 acks=all:选对 ISR 才是关键

acks=all也不是无敌的,很多人的理解是“设成 all 就绝对不丢了”,这其实是个误区。acks=all的含义是 Leader 必须等待当前 ISR(In-Sync Replica,同步中副本)集合里的所有副本都写入成功,才返回成功。注意,ISR 集合是动态变化的。

如果一个 Follower 跟不上水位,它会被踢出 ISR,那么acks=all实际上只需要等 ISR 里剩下的那些副本确认。极端情况下,ISR 里只有 Leader 自己,这时候acks=all的可靠性就退化成acks=1了。

所以我有一个习惯性建议:把min.insync.replicas这个参数设置好,比如设置成 2。它的作用是最少要有多少个 ISR 副本在线,否则分区直接不可写。这样配合acks=all,可靠性才能真正兜住底。不要只改一个参数,这俩是配套的,少了谁都有漏洞。

2.4 异步发送、缓冲区和超时:三个隐藏的地雷

除了 acks,生产端还有三个隐藏炸弹:异步发送、max.block.ms、buffer.memory。

很多人写代码时喜欢用异步发送,因为同步发送在 Broker 繁忙时会有明显的 RT 上涨。异步发送本身没问题,问题在于很多人发了消息之后没有正确注册回调,或者回调里只是打了行日志就没后续。消息发送失败后,你毫无感知,这对“丢没丢消息”的判断是致命的。

缓冲区参数也一样。如果生产者的buffer.memory设置太小,或者max.block.ms设置得过于保守,在业务洪峰来临时,KafkaProducer 会直接抛出异常。如果你没做重试和补偿,这条消息就静默丢失了。建议把异常处理落到实处,发送失败至少要有重试机制,而且重试时必须关注enable.idempotence(幂等生产者),避免重试造成消息重复。

3. 服务端:Broker 层面的丢失,其实最容易被忽视

3.1 刷盘策略:断电的一瞬间,内存里的都白写

Broker 端的第一个“丢消息”隐患就是刷盘策略。Kafka 利用操作系统 PageCache 来提升写入性能,生产数据会先写入 PageCache,再由系统异步刷到磁盘。这个机制在绝大多数情况下都没问题,性能高、吞吐大。

但一旦发生断电或者内核崩溃,PageCache 里还没有落盘的数据就会全部丢失。有些团队会调整log.flush.interval.messages和log.flush.interval.ms来强制更频繁地刷盘,但这两个参数和吞吐量是冲突的。你不能又想要极致吞吐又想要断电不丢数据,这不现实。

我实际遇到的案例里,很多公司压根没注意日志目录的磁盘是不是 SSD、IO 是否稳定。磁盘性能差会导致 Follower 复制延迟升高,进而被踢出 ISR,最终影响整个分区的可用性。这里我想跟大家说的经验是:不要把 Kafka 部署在性能很差的云盘上,否则你会遇到各种诡异的消息延迟和副本同步问题。

3.2 副本数是 2?丢一半的概率可不算小

很多人创建主题时只设置了 1 个副本,或者图省钱只设 2 个副本。副本是 1 的时候,机器一挂,分区数据全没,这种场景已经不能叫“丢消息”了,应该叫“数据毁灭”。副本数设置为 2,其实也并没有好到哪里去,因为它只能容忍一台机器故障,遇到机房级别的断电,依然可能同时挂掉两个副本。

我自己在给项目做容量规划的时候,一般建议核心业务主题至少 3 个副本。如果你说我资源紧张,至少保证min.insync.replicas=2,并且磁盘、机器尽量分散到不同的故障域。记住一个核心原则:可靠性是设计出来的,不是运气碰出来的。你可以用最小副本数跑测试环境,但生产环境请别抱侥幸心理。

另外想说一下,副本数不是设置完就完事的。Kafka 的副本分布和机器的机架信息有关,如果所有副本恰好都分配在同一台物理机上,那等于没有副本。你可以通过kafka-reassign-partitions.sh来做副本重分配,至少看一眼分区副本的分布情况再上线。

3.3 Leader 选举与 ISR 收缩:你以为的高可用是暂时的

一条消息能不能被安全消费,其实取决于它是否完成了“高水位”同步。高水位就是所有 ISR 副本都同步到的位置。如果你的消费者只读取了 Leader 上的数据,而一条消息还没来得及同步到所有副本,Leader 挂掉时,新选举出来的 Leader 上可能没有这条消息,消费端就会觉得消息丢了。

还有一个很多初学者会忽视的点:Leader 切换时,消息的偏移量可能会重置。比如旧 Leader 上有 offset 是 0~100 的数据,但新 Leader 只同步到了 offset 80,那么旧 Leader 上 80~100 的数据对客户端来说就是“被丢了”。严格意义上说,这不是 Kafka 丢了你的消息,而是你读取的是一个并不完整的副本数据。

所以我们要依赖min.insync.replicas和 ISR 机制,把“数据未完整同步前不可用”变成硬约束。你可以把 ISR 收缩、Leader 切换的日志全部打开来观察,一旦发现 ISR 频繁变化,就要注意了,这往往意味着某台机器磁盘 IO 很高、GC 停顿严重或者网络有问题。表面上是在“丢消息”,实质上是在“丢副本”。

3.4 日志清理和配额:别让物理资源替你“丢消息”

还有一个冷门但真实会发生的场景:日志清理。Kafka 的主题默认开启了日志删除机制,超过retention.ms或者超过retention.bytes的旧数据会被自动清理。如果你消费端挂了好几天,回来继续消费的时候发现从头开始读已经读不到多少数据了,那不是 Kafka 丢消息,而是消息被过期清理了。

这种问题常发生在测试环境或者运维不规范的团队里。消费者离线太久,清理线程把旧日志都删掉了。我建议你根据数据保留需求设置合理的retention.ms,同时监控消费组 Lag 的监控指标,发现 Lag 持续增长就要及时处理,别等消息被清理了再去找原因。

还有一种情况是磁盘空间被打满,导致日志写入失败,Kafka 会直接停止接受写入。这个场景比较像“丢消息”,其实是集群进入保护状态。日常运维中,必须为日志目录预留充足的空间,并配置好磁盘使用率告警,别让物理资源掐住了消息通道的脖子。

4. 消费端:处理不掉的“丢”,其实是消费位点的问题

4.1 自动提交位点:程序一崩,一批消息直接“丢了”

消费端的“丢消息”往往最具有迷惑性,因为它在很多时候表现为“消费者明明拉取了消息,结果却像没有消费过一样”。最经典的场景就是开启enable.auto.commit=true,默认每 5 秒自动提交一次消费位点。

想象一个场景:消费者把一批消息拉了下来,业务处理过程中 Redis 超时,代码抛了异常,于是这批消息还没有处理完,但 Kafka 客户端的后台线程已经按照落地时间提交了 offset。等重启消费者之后,Kafka 会认为这批消息已经消费完毕,直接跳到后面去了。你的业务数据缺失,看起来就像 Kafka 把消息丢了。

最让我无语的是,有些团队把错误处理逻辑做成打日志然后继续消费,这样连“消费失败”的痕迹都看不到。下游数据对不上账的时候,翻业务日志全是异常后跳过继续往下走的记录。这里我给大家一个最基础也最有效的建议:核心业务里不要自动提交位点,手动提交,并且做到“先处理完业务,再提交 offset”。

4.2 enable.auto.commit=false 真的就不丢了吗

手动提交位点解决了“处理一半就提交”的问题,但如果你用手动提交的方式,又容易引发另一个常见问题:位点提交失败。比如你的消费者处理完一条消息,准备提交 offset 的时候,刚好发生 Rebalance,这时第一次提交失败。如果你没做重试,重启之后就会从这个旧位点重新消费。

这种情况和“丢消息”相反,它会导致“重复消费”,而不是消息丢失。但引出的另一个极端是,有些人为了追求效率,在业务处理之前就提前提交 offset,这样下一次消费者挂掉时,这些消息就永远不会再被读到。这不叫丢,但在业务效果上跟丢了一模一样。

所以我特别想强调一个观点:在 Kafka 的消费模型里,“不丢消息”和“不重复消息”这俩目标本质上就是冲突的。你只能选择“至少一次”或者“恰好一次”,没有完美的中间态。如果业务要求绝对不能丢,那你就要接受可能会有重复,并且在消费端做幂等。比如用 Redis 记录消息唯一 ID,或者直接把幂等判断做到数据库里。

4.3 重平衡:踢你出群的那一刻,消息去哪里了

重平衡(Rebalance)是消费端丢消息的另一个重灾区。Kafka 消费者通过消费组协调分区分配,当消费者数量变化、消费组订阅关系变化,或者消费者心跳超时的时候,会触发重平衡。重平衡本身不是问题,问题在于重平衡发生时,未提交位点的分区会被重新分配给其他消费者。

打个比方,分区 0 之前分配给了消费者 A,A 拉了几百条消息正在处理,还没提交位点,结果网络抖动,消费组触发了重平衡,分区 0 分给了消费者 B。B 会从上次提交的位点开始消费,那些 A 拉下来但没来得及落库的数据就“丢”了。

当然现代 Kafka 的协同式重平衡已经比老版本温和了不少,但老版本的问题在新版本里依然可能因为心跳超时、session.timeout.ms设置不合理而出现。我给团队的建议是,不要把max.poll.records调得过大,单次拉取太多消息处理时间太长,容易导致消费者超过max.poll.interval.ms被判定为不可用,从而触发重平衡。

4.4 从“至少一次”到“恰好一次”:别把丢和重混为一谈

最后一条是关于语义边界的认知。很多人说“我想做到不丢消息”,其实他们内心的真实诉求是“我一条消息都不多、一条都不少”。要做到一条都不少,需要引入 Kafka 的幂等生产者和事务机制,配合消费端的幂等消费,也就是端到端恰好一次。

但这里有一个现实问题:绝大多数业务系统追求的是“至少一次”加上消费端的业务幂等。比如你用订单号做主键做 upsert,重复插入和插入一次效果一样,你就不怕重复消费。反而是一味追求事务机制,会带来比较大的性能损耗和实现复杂度。站在我个人的角度,除非你是做支付、对账一类强一致系统,否则优先采用“生产幂等 + 消费幂等”这个组合,性价比最高。

搞清楚语义边界,你才不会被“丢消息”和“重复消息”这两个词来回绕晕。排查问题的时候,要搞清楚用户说的“丢”,到底是真的丢了,还是重复消费中被掩盖了,还是只是延迟变高导致看起来像是没消息。这一步清晰了,很多时候你根本不需要调一堆参数。

5. 排查实录:遇到丢消息,我是怎么一步步查的

5.1 第一步:先分清楚是“没生产成功”还是“没消费到”

排查丢消息,最忌讳一上来就翻 Kafka 源码或者乱调参数。我自己的固定套路是先分域:到底是生产端没发出去,Broker 端没存住,还是消费端没消费掉。

生产端好判断,你只要看生产者日志里有没有异常,回调里有没有记录失败消息。如果回调里记录了大量失败,那问题就定位在生产端,先看max.block.ms是否触发、重试次数是否不够、目标 topic 是否存在。如果回调都是空的,可以去 Broker 端看这个 topic 的落盘总消息数。

消费端则要看消费组的 Lag。用命令行工具执行kafka-consumer-groups.sh --describe --group xxx,观察每条分区的CURRENT-OFFSET和LOG-END-OFFSET。如果 Lag 很高,说明消息没被及时消费完,先排查消费者是否有异常退出。如果 Lag 是 0 但下游数据缺失,那大概率是“处理失败但提交成功”,需要回源日志看业务处理状态。

5.2 第二步:用数据验证,别靠感觉

很多时候大家判断丢消息都是“感觉丢了”。比如下游说某个时段的数据少了,然后就开始抓狂。我建议遇到这种场景先做个端到端的数字验证:生产者端统计发送成功总数,Broker 端统计某个 topic 的累积消息量,消费者端统计处理成功数。

三个数字一对,问题到底出在哪一段就清楚了。如果生产者发送成功数比 Broker 的入站消息数多,那是生产者回调处理有问题,消息实际上没发到 Broker。如果 Broker 端有但消费者端没有,那就是消费端处理逻辑或提交位点的问题。如果 Broker 端本身就少了,再继续查副本、磁盘和 ISR 的情况。

这里分享一个我常用的手法:在生产端往消息里注入唯一 ID,消费端用这个 ID 做去重和统计,这样几乎可以把每条消息的流转轨迹都捞出来。日志长期留存的情况下,排查“丢消息”会变得非常直观。不要嫌麻烦,这样的投入在关键时刻能救你一命。

5.3 第三步:常见“丢消息”场景速查表

下面是我这些年排查中遇到最常见的几类“丢消息”场景,直接做成速查表,方便你对照自己的系统:

场景表现潜在原因核心检查点
生产端偶发超时,消息缺失网络抖动、delivery.timeout.ms过小增加重试次数,调大超时时间
极端情况下 Leader 切换,少量消息丢失acks=1或acks=all但min.insync.replicas太小设置acks=all,并配套min.insync.replicas=2
磁盘满或 IO 异常,写入失败机器磁盘容量/性能不足监控磁盘,配置告警,增加节点
消费者重启后部分数据缺失自动提交位点,处理失败但已提交关闭自动提交,先处理再提交
消费组重平衡后数据缺失未提交位点 + Rebalance 触发重分配检查心跳超时、max.poll.interval.ms
消息还在但消费不到日志过期被清理设置合理retention.ms,监控 Lag

这张表不是标准答案,但它基本覆盖了日常 90% 的丢消息故障。如果你遇到的问题不在这张表里,那再去翻kafka.server:type=KafkaRequestHandlerPool、RequestQueueTimeMs这些监控指标,看有没有请求积压或者处理瓶颈。

5.4 几条我踩过坑之后留下的经验

第一条经验是:永远不要依赖环境默认值。Kafka 的默认配置是“均衡”的,但用在你的业务上未必合适。每个团队都要根据自己的消费频率、允许延迟、数据重要性去调整配置。最好把关键参数固化到你们的部署模板里,新项目直接套用。

第二条经验是:监控一定做在前头。Kafka 的指标非常多,我不建议一开始就堆一堆 Dashboard,先盯几个关键项就够了,生产端成功消息数、消费组 Lag 值、ISR 副本数是否健康、Broker 磁盘剩余空间。这四个指标能覆盖大多数丢消息预兆。

第三个比较个人化的体会是:不要在工作中试图制造完美的“绝不丢失”系统。我见过不少同事为了追求极端可靠性,设计出了一套非常复杂的事务方案,最终运行维护成本剧增,反而是引入了更多故障点。成熟的架构会接受“去重”而不是“杜绝重复”,用幂等去消化不确定性,而不是把一切都堵死在消息系统这层。

6. 附:值得长期关注的 Kafka 可靠性配置清单

如果你看到这里,还不知道从哪里入手,我直接给你一份我常用的配置清单,里面每一项都是围绕“降低丢消息概率”来设计的。

主题级别配置:

  • replication.factor=3,至少也要 3 个副本,核心主题别低于这个数。
  • min.insync.replicas=2,保证至少有两个副本同步完成才可写。
  • retention.ms要匹配消费端的最大离线容忍时间,别设太短。
  • cleanup.policy=delete保留默认即可,除非你需要 compact。

生产者级别配置:

  • acks=all。
  • enable.idempotence=true,让生产者自动处理重复与乱序。
  • retries设置一个较大值,比如Integer.MAX_VALUE,配合delivery.timeout.ms控制总时长。
  • max.in.flight.requests.per.connection在开启幂等后可以保持默认 5,不用刻意调为 1。
  • buffer.memory根据你的单机峰值发送量调大,别默认 32MB,高峰期很容易满。
  • max.block.ms建议在 3000 以上,不然一阻塞就抛异常。

消费者级别配置:

  • enable.auto.commit=false,手动提交异步入库成功后再提交。
  • auto.offset.reset根据业务选,默认latest对新消费组来说会直接跳过已有消息,如果你需要“从最早开始消费”就设成earliest。
  • max.poll.interval.ms要大于你的单次批量处理耗时余量,建议至少 3 倍余量。
  • session.timeout.ms不要太高也不要太低,默认 45 秒在很多场景够用,但要注意网络抖动时太短容易频繁重平衡。

这些参数不是一成不变的,但你可以把它们当成一个比较可靠的基线。业务跑起来之后,再根据监控数据微调单个参数,每次修改都要能说清为什么改,不要凭感觉拍脑袋。我自己就是这样从一堆“丢消息”的事故里熬出来的,希望对你有帮助。

返回列表