Kafka 在很长一段时间里都被当成消息中间件里的“天花板”来用吞吐量高、可扩展性好、还能做存储。但凡是线上实际用过 Kafka 的团队几乎都经历过一次“到底丢没丢消息”的争论。有人说是 Kafka 的问题有人说是自己代码的问题还有人说是网络抖动导致的偶发故障。实际上我做了这么多年 Kafka 相关的项目结论非常明确Kafka 本身很少会“无理由丢消息”绝大多数丢消息的根源都出在生产者配置、消费者提交位点的时机以及你对确认机制的理解上。这篇文章我会从生产端、Broker 端、消费端三个链路分别拆解把“丢消息”这件事彻底讲透。不管你是在做日志采集、实时数仓还是普通的后端消息推送这篇内容都值得对照你当前的配置逐条查一遍。你会发现很多丢消息不是玄学而是配置和代码写得不够谨慎。1. 先说结论丢消息从来都不是 Kafka 一个环节的问题1.1 我对“Kafka 丢消息”这句话的理解很多人一搜“Kafka 丢消息”第一反应是 Kafka 这个中间件本身有没有 bug。其实 Kafka 作为分布式系统单从设计上讲它已经尽最大努力去保证消息不丢了。问题在于这个“保证”是有前提条件的比如生产者要设置合理的确认机制、Broker 端要有足够的副本数、消费者要自己管理好偏移量。你如果把保证可靠性的责任全部甩给 Kafka自己代码里毫不在意地往生产端一塞、消费端一拉那丢消息几乎是必然的。所以我的第一句话永远是Kafka 不丢消息是有条件的你要做的不是问 Kafka 为什么丢而是问自己在哪个环节破坏了它的可靠性边界。1.2 先想清楚消息从进到出要经过哪几道门一条消息从业务系统产生到最终被下游消费者处理完整的路径是生产者把消息发送给 Kafka BrokerBroker 把消息写入分区对应的日志文件并完成副本同步消费者从 Broker 拉取消息处理业务逻辑消费者提交偏移量offset确认这条消息已经被消费。这里面每一道门都有可能把消息弄丢。生产端可能因为发送失败、回调没处理、缓冲区满了而丢Broker 端可能因为刷盘不及时、副本同步失败、Leader 切换而丢消费端可能因为位点提交太快、程序崩溃、重平衡而丢。你光盯着某一端排查是没用的。我见过不少团队 debug 了好几天最后发现是消费端业务异常导致位点已经提交了但消息没处理完成这种“假成功”最容易让人误以为是 Kafka 在丢消息。所以先建立这个全局链路认知后面排查起来才有方向。2. 生产端你还没发出去消息就已经丢了2.1 acks0 的代价说发就发后果自负生产者端第一个经典坑就是 acks 配置。如果你把acks设成 0意味着生产者不需要等待 Broker 的任何确认消息发出去就算完事。这个配置下吞吐量最高但代价就是只要网络发生抖动、请求还没到 Broker或者 Broker 写入过程中出现问题生产者这边完全无感知。我实际帮团队排查过一个采集系统的丢数据问题客户端采集到日志后直接往 Kafka 里怼配置里acks0结果高峰期数据丢失率在 0.5% 左右看着不高但对于需要精准对账的业务来说完全不可接受。后来把acks改成 1丢失率立刻降下来了。在这个配置下你所谓“发消息成功”只是客户端把数据交给了 socket 缓冲区并不是真的写到了 Kafka。所以生产环境里除了极少数允许丢失的监控指标类数据我认为一般不要用acks0。真到了丢不起的场合再回来的查证成本远高于那点吞吐量收益。2.2 acks1 的中间地带Leader 挂了你就傻眼acks1是很多团队里的默认选择意思是生产者只需要等 Leader 副本写入成功就返回成功。这个配置在实际吞吐量和可靠性之间确实比较均衡很多业务也都跑在这种模式下。但你要清楚它的死角如果 Leader 在写入完成后、还没等其他副本同步的时候挂了这条消息对消费者来说就是“丢失”的。举个例子你给主题配置了一个分区三个副本某台机器上的 Leader 副本把消息写盘后数据还没同步到其他两个 Follower结果这台机器突然断电。Kafka 会从剩下的 Follower 里选一个新的 Leader之前那批没有同步完成的消息就永远没了。在这个环节你要是只盯着生产者有没有收到成功回调那是不够的。acks1只保证“那一刻的 Leader 写入成功”不保证这条消息在集群里永久存活。对于账务类、订单类这种强一致场景我建议直接上acksall别太抠那一点点延迟。2.3 最稳妥的 acksall选对 ISR 才是关键acksall也不是无敌的很多人的理解是“设成 all 就绝对不丢了”这其实是个误区。acksall的含义是 Leader 必须等待当前 ISRIn-Sync Replica同步中副本集合里的所有副本都写入成功才返回成功。注意ISR 集合是动态变化的。如果一个 Follower 跟不上水位它会被踢出 ISR那么acksall实际上只需要等 ISR 里剩下的那些副本确认。极端情况下ISR 里只有 Leader 自己这时候acksall的可靠性就退化成acks1了。所以我有一个习惯性建议把min.insync.replicas这个参数设置好比如设置成 2。它的作用是最少要有多少个 ISR 副本在线否则分区直接不可写。这样配合acksall可靠性才能真正兜住底。不要只改一个参数这俩是配套的少了谁都有漏洞。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.replicas2并且磁盘、机器尽量分散到不同的故障域。记住一个核心原则可靠性是设计出来的不是运气碰出来的。你可以用最小副本数跑测试环境但生产环境请别抱侥幸心理。另外想说一下副本数不是设置完就完事的。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.committrue默认每 5 秒自动提交一次消费位点。想象一个场景消费者把一批消息拉了下来业务处理过程中 Redis 超时代码抛了异常于是这批消息还没有处理完但 Kafka 客户端的后台线程已经按照落地时间提交了 offset。等重启消费者之后Kafka 会认为这批消息已经消费完毕直接跳到后面去了。你的业务数据缺失看起来就像 Kafka 把消息丢了。最让我无语的是有些团队把错误处理逻辑做成打日志然后继续消费这样连“消费失败”的痕迹都看不到。下游数据对不上账的时候翻业务日志全是异常后跳过继续往下走的记录。这里我给大家一个最基础也最有效的建议核心业务里不要自动提交位点手动提交并且做到“先处理完业务再提交 offset”。4.2 enable.auto.commitfalse 真的就不丢了吗手动提交位点解决了“处理一半就提交”的问题但如果你用手动提交的方式又容易引发另一个常见问题位点提交失败。比如你的消费者处理完一条消息准备提交 offset 的时候刚好发生 Rebalance这时第一次提交失败。如果你没做重试重启之后就会从这个旧位点重新消费。这种情况和“丢消息”相反它会导致“重复消费”而不是消息丢失。但引出的另一个极端是有些人为了追求效率在业务处理之前就提前提交 offset这样下一次消费者挂掉时这些消息就永远不会再被读到。这不叫丢但在业务效果上跟丢了一模一样。所以我特别想强调一个观点在 Kafka 的消费模型里“不丢消息”和“不重复消息”这俩目标本质上就是冲突的。你只能选择“至少一次”或者“恰好一次”没有完美的中间态。如果业务要求绝对不能丢那你就要接受可能会有重复并且在消费端做幂等。比如用 Redis 记录消息唯一 ID或者直接把幂等判断做到数据库里。4.3 重平衡踢你出群的那一刻消息去哪里了重平衡Rebalance是消费端丢消息的另一个重灾区。Kafka 消费者通过消费组协调分区分配当消费者数量变化、消费组订阅关系变化或者消费者心跳超时的时候会触发重平衡。重平衡本身不是问题问题在于重平衡发生时未提交位点的分区会被重新分配给其他消费者。打个比方分区 0 之前分配给了消费者 AA 拉了几百条消息正在处理还没提交位点结果网络抖动消费组触发了重平衡分区 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 切换少量消息丢失acks1或acksall但min.insync.replicas太小设置acksall并配套min.insync.replicas2磁盘满或 IO 异常写入失败机器磁盘容量/性能不足监控磁盘配置告警增加节点消费者重启后部分数据缺失自动提交位点处理失败但已提交关闭自动提交先处理再提交消费组重平衡后数据缺失未提交位点 Rebalance 触发重分配检查心跳超时、max.poll.interval.ms消息还在但消费不到日志过期被清理设置合理retention.ms监控 Lag这张表不是标准答案但它基本覆盖了日常 90% 的丢消息故障。如果你遇到的问题不在这张表里那再去翻kafka.server:typeKafkaRequestHandlerPool、RequestQueueTimeMs这些监控指标看有没有请求积压或者处理瓶颈。5.4 几条我踩过坑之后留下的经验第一条经验是永远不要依赖环境默认值。Kafka 的默认配置是“均衡”的但用在你的业务上未必合适。每个团队都要根据自己的消费频率、允许延迟、数据重要性去调整配置。最好把关键参数固化到你们的部署模板里新项目直接套用。第二条经验是监控一定做在前头。Kafka 的指标非常多我不建议一开始就堆一堆 Dashboard先盯几个关键项就够了生产端成功消息数、消费组 Lag 值、ISR 副本数是否健康、Broker 磁盘剩余空间。这四个指标能覆盖大多数丢消息预兆。第三个比较个人化的体会是不要在工作中试图制造完美的“绝不丢失”系统。我见过不少同事为了追求极端可靠性设计出了一套非常复杂的事务方案最终运行维护成本剧增反而是引入了更多故障点。成熟的架构会接受“去重”而不是“杜绝重复”用幂等去消化不确定性而不是把一切都堵死在消息系统这层。6. 附值得长期关注的 Kafka 可靠性配置清单如果你看到这里还不知道从哪里入手我直接给你一份我常用的配置清单里面每一项都是围绕“降低丢消息概率”来设计的。主题级别配置replication.factor3至少也要 3 个副本核心主题别低于这个数。min.insync.replicas2保证至少有两个副本同步完成才可写。retention.ms要匹配消费端的最大离线容忍时间别设太短。cleanup.policydelete保留默认即可除非你需要 compact。生产者级别配置acksall。enable.idempotencetrue让生产者自动处理重复与乱序。retries设置一个较大值比如Integer.MAX_VALUE配合delivery.timeout.ms控制总时长。max.in.flight.requests.per.connection在开启幂等后可以保持默认 5不用刻意调为 1。buffer.memory根据你的单机峰值发送量调大别默认 32MB高峰期很容易满。max.block.ms建议在 3000 以上不然一阻塞就抛异常。消费者级别配置enable.auto.commitfalse手动提交异步入库成功后再提交。auto.offset.reset根据业务选默认latest对新消费组来说会直接跳过已有消息如果你需要“从最早开始消费”就设成earliest。max.poll.interval.ms要大于你的单次批量处理耗时余量建议至少 3 倍余量。session.timeout.ms不要太高也不要太低默认 45 秒在很多场景够用但要注意网络抖动时太短容易频繁重平衡。这些参数不是一成不变的但你可以把它们当成一个比较可靠的基线。业务跑起来之后再根据监控数据微调单个参数每次修改都要能说清为什么改不要凭感觉拍脑袋。我自己就是这样从一堆“丢消息”的事故里熬出来的希望对你有帮助。