你打开监控看板发现昨天某个业务的消息总量对不上少了三千条。查日志、查消费者、查broker折腾一晚上最后同事轻描淡写说了一句“哦那个topic的acks配置没改生产者重启的时候掉了几台。”——这种场景在Kafka技术群里几乎每天都在上演。Kafka丢消息表面上是某个配置项的问题但往深了说往往牵涉生产端、broker端、消费端三个环节还有你根本没意识到的操作系统与硬件因素。这篇文章不聊空泛原理我直接把Kafka丢消息的根源拆成几个看得见、摸得着的检查点把生产环境里真正踩过的坑和排查套路完整摆出来。无论你是刚搭完3节点集群的运维新人还是被消息延迟搞到焦头烂额的开发里面提到的每一条我都建议你对照自己的环境过一遍。1. 丢消息的根源三个环节逐个排查很多人一听到“Kafka丢消息”第一反应就是broker没把数据落盘。但根据我这一线排查经验绝大多数丢消息都不是broker一个人的锅而是生产端、broker端、消费端三个环节里的某一个粗心配置导致的。咱们按消息流向从源头到末尾挨个拆。1.1 生产端丢消息acks和retries没有你想象的那么简单生产端丢消息最常见的原因就是acks参数设置不当。Kafka生产者发送消息时acks有三个取值0、1、all。acks0生产者发完即焚不等待broker任何确认。这个配置下丢消息是必然的只是时间问题。acks1只要leader副本写入成功就返回确认。leader所在节点宕机、数据还没来得及同步给follower这部分消息就丢了。acksall消息被所有ISR同步副本都确认后才返回成功可靠性最高但吞吐会有所下降。我看到很多团队在压测时为了冲吞吐把acks调成0或1上线后忘了改回来这是最典型的丢消息姿势。生产环境里请直接使用acksall。另外还有一个隐蔽的坑retries。默认情况下retries不为0但对“发送失败是否需要重试”这件事很多开发者会忽略max.in.flight.requests.per.connection的影响。如果这个值大于1且开了重试消息的发送顺序可能发生颠倒。更稳妥的做法是开启幂等能力同时保证不丢不重Properties props new Properties(); props.put(ProducerConfig.ACKS_CONFIG, all); props.put(ProducerConfig.RETRIES_CONFIG, Integer.MAX_VALUE); props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5); props.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, 120000);这里有个概念必须说透幂等只能保证“单分区内、单会话内”不重不丢。如果生产者重启或者broker端发生了leader切换部分跨会话的消息仍然存在重复风险。所以别把幂等当万能药该做的业务幂等校验比如用消息主键去重省不了。生产端还有一个容易被漏掉的点future.get()有没有被真正执行。很多人在用KafkaProducer发送时只调用send()却从不检查返回的Future或者把异常吞掉。发送失败、序列化失败、broker不可达时回调不处理消息就安静地消失了。我的建议是统一封装发送函数同步等待结果或注册回调至少保证异常能被感知。1.2 Broker端丢消息副本数够为什么还会丢broker端丢消息通常可以归纳为三种情况副本数不够、ISR机制没配置好、底层存储出了问题。先说说副本数。很多小团队为了省资源topic的replication.factor设置成1也就是只有一个副本。这种情况下一旦leader节点宕机、磁盘损坏消息就跟着没了。3节点集群起步已经是行业共识topic级别最少也要replication.factor3并且把min.insync.replicas设置为2。这样才能保证leader写入后至少还有一个follower同步完成消息才算“成功”。但光有副本还不行。Kafka的同步机制里ISR是动态收缩的。如果一个follower长时间不赶上leader由replica.lag.time.max.ms控制它会被踢出ISR。如果所有follower都不在ISR里只有leader自己而你又设置了min.insync.replicas2那么生产请求会直接报NotEnoughReplicasException。这个异常看似是“写入失败”其实是Kafka在保护数据可靠性——它宁愿拒绝写入也不让你把消息写进一个完全没有冗余的副本里。这里有个反直觉的坑unclean.leader.election.enable。默认值是false意思是当leader挂了ISR里没有可用副本时不允许“非ISR副本”竞选leader。如果你手贱把它改成trueKafka会允许一个落后的副本成为新leader那么它缺失的那些消息就彻底“没”了。这是我看过最危险的一个开关建议保持false。broker端还有一个很少人关注但确实会导致数据丢失的物理因素磁盘。Kafka重度依赖顺序写但不管顺序写多快最终数据还是要落到磁盘上。如果log.dirs配置成了系统盘操作系统日志、swap、其他应用IO一拥而上磁盘写入抖动极端情况下会触发log segment异常或者broker进程崩溃。数据落在page cache里没刷盘节点一断电没来得及刷盘的消息就会丢失。所以生产环境一定要用独立的SSD或NVMe磁盘并且把log.flush.interval.messages、log.flush.interval.ms调成适合业务的持久化策略。虽然Kafka默认依靠操作系统刷盘但对于金融、订单这类“一条都不能丢”的场景建议显式设置较低的门槛用一点吞吐换可靠性。1.3 消费端丢消息大部分人都栽在offset上消费端丢消息最典型的场景就是enable.auto.committrue。默认情况下消费者每隔auto.commit.interval.ms默认5000ms自动提交位移。如果你的业务逻辑是先提交位移、后处理消息或者异步处理消息尚未完成消费者进程崩溃或触发rebalance那些“已经提交offset但没处理完”的消息就永远不会被再次消费相当于丢了。这种丢消息和broker无关纯粹是消费语义问题。生产环境我强烈建议改成props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);然后手动提交。提交时机怎么选处理完业务逻辑之后、再提交位移。例如while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecordString, String record : records) { handleBusiness(record); } consumer.commitSync(); }这里要注意commitSync和commitAsync的区别。commitAsync速度快但提交失败时不会重试因为可能出现“后提交成功、前提交失败”导致位移乱序。我建议业务关键路径上用commitSync或者两者配合先异步提交、失败时同步补偿提交。还有一类“看似不丢、实则重复”的情况严格说不是丢消息而是重复消费。业务侧如果不做幂等重复消费同样会引发数据混乱。所以我说消费端做得好的团队一定是“手动提交业务幂等”两手抓。2. 选型阶段的隐患Kafka、RabbitMQ、RocketMQ怎么选才不丢消息很多人在项目起步阶段就选了Kafka没想过这套选型本身也在影响“丢消息”这件事。我在选型实战对比里踩过不少坑这里说点真实的体会。2.1 三种消息队列的存储模型差异我习惯把三者用一句话概括Kafka是“日志流”模型RabbitMQ是“队列”模型RocketMQ是两者之间的“合并改进型”。Kafka把消息按partition追加写入log文件消费者通过offset定位消息。这个设计赋予了Kafka超高的吞吐因为它是顺序写磁盘、批量发送、零拷贝读取硬件利用效率极高。但可靠性模型里它依赖副本同步和ISR机制配置错了就是丢消息重灾区。RabbitMQ是真正的队列模型。消息进入exchange按路由键分发到queue消费者消费后要显式ackbroker才删除消息。RabbitMQ在no_ackfalse的默认语义下每条消息的生命周期都很清晰丢消息的概率天然低很多但它做不到Kafka那种超高吞吐适合事务性场景。RocketMQ的存储模型是“commitlog consumequeue”倒有点像Kafka和RabbitMQ的折中。它同时支持事务消息、延迟消息和消费端的重试队列尤其在“消息发送成功但消费失败”的场景下有broker端的重试机制可靠性设计比Kafka更直白。2.2 可靠性与吞吐的权衡如果单论“丢消息的概率”在默认配置下RabbitMQ和RocketMQ都比Kafka更“稳”因为它们把消息确认和重试机制做得更靠前。Kafka的可靠性不是不行而是需要你手动调优的地方太多了。举一个具体例子RabbitMQ中一条消息如果消费者没有ack它会一直待在队列里即使消费者挂了消息也不会丢。Kafka则不同如果消费者进程在处理消息之前就提交了offset这条消息就再也找不回来了。这种设计差异导致很多人从RabbitMQ转到Kafka时会很不适应。换一个角度看性能Kafka单分区顺序写可以达到数十万甚至上百万TPSRabbitMQ单队列吞吐受限于ack确认时的往返延迟通常在万级到十万级。RocketMQ则在吞吐和延迟之间做了平衡但运维复杂度比Kafka还高不是所有团队都扛得住三件套NameServer、Broker、Dashboard的维护成本。2.3 选型决策清单结合真实项目经验我建议这样选日志采集、用户行为追踪、数据管道、事件流处理无脑选Kafka。业务订单、支付通知、短信服务、需要死信队列和灵活ack语义选RabbitMQ更省心。需要事务消息、定时消息、而且团队有Java基础沉淀RocketMQ是可靠性与吞吐的折中方案。有人问我“我就是要用Kafka但又怕丢消息怎么办”我的回答是那就用Kafka同时把第1节和第3节提到的配置全部做对。Kafka不适合对可靠性的要求极高、又没人愿意维护参数细节的小团队。选型不丢人选错了硬扛才丢人。3. 部署与调优3节点集群里最容易踩的坑Kafka集群部署本身不复杂复杂的是部署完之后那些看起来“没问题”的默认参数。我以一个3节点集群为模板把每一步的关键配置和背后的原因说清楚。3.1 三节点集群的规划与安装3个节点是Kafka集群的最小高可用形态。节点的broker.id必须唯一我用的是0、1、2。每个节点的listeners和advertised.listeners一定要分开配置尤其是服务器有多块网卡、前后端分离网络的情况下。advertised.listeners配置不当外部客户端会拿到一个无法连接的地址线上线下排查起来极其痛苦。# server.properties 核心配置 broker.id0 listenersPLAINTEXT://0.0.0.0:9092 advertised.listenersPLAINTEXT://192.168.1.10:9092 log.dirs/data/kafka-logs num.partitions3 default.replication.factor3 min.insync.replicas2 replica.lag.time.max.ms30000 unclean.leader.election.enablefalse auto.create.topics.enablefalse安装时两个坑我必须提一下。一个是JDK版本与Kafka版本的兼容性。老旧的JDK8配新版Kafka或者反过来都可能在运行一段时间后出现奇怪的网络异常。另一个是系统文件句柄数。Kafka broker持有大量日志文件句柄生产环境必须把ulimit -n调到65535以上否则运行几天后就会出现“Too many open files”直接服务不可用。Topic配置同样重要。很多人只改了default.replication.factor却忘了已有topic的副本数不会自动更新。3节点集群里新建topic建议显式声明kafka-topics.sh --bootstrap-server 172.16.0.1:9092,172.16.0.2:9092,172.16.0.3:9092 \ --create --topic your-topic \ --partitions 6 --replication-factor 3 \ --config min.insync.replicas2分区数这里多说一句不要为了吞吐无脑把分区数设得很高。分区数过多会带来broker端文件句柄增长、leader切换耗时增加、消费者rebalance时间变长等问题。分区数的合理范围通常建议是broker总核数到总核数×2之间。3.2 组关键参数备查我把生产环境里影响可靠性的一组参数整理成一张表大家直接对着改参数推荐值作用acksall要求ISR全部确认min.insync.replicas2最少同步副本数replication.factor3副本总数量unclean.leader.election.enablefalse禁止非ISR副本竞选leaderenable.auto.commitfalse关闭自动位移提交replica.lag.time.max.ms30000控制ISR收缩时间log.flush.interval.messages10000消息数刷盘阈值log.flush.interval.ms500时间刷盘阈值注意最后两个刷盘参数Kafka的默认行为是交给操作系统来刷盘理论上进程崩溃不丢、节点断电可能丢。如果要求“断电也不丢”必须显式调低刷盘阈值代价是吞吐明显下降。3.3 用AdminClient和可视化工具确认集群状态靠肉眼去记集群配置迟早翻车我强烈建议用AdminClient写一个小工具把集群状态和topic配置拉出来看。下面这段代码可以输出每个topic的副本分配和ISR状态Properties props new Properties(); props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, 172.16.0.1:9092,172.16.0.2:9092,172.16.0.3:9092); try (AdminClient admin AdminClient.create(props)) { DescribeTopicsResult describeTopics admin.describeTopics(List.of(your-topic)); describeTopics.all().get().forEach((topic, detail) - { System.out.println(topic: topic); detail.partitions().forEach(p - { System.out.println( partition p.partition() leader p.leader() replicas p.replicas() isr p.isr()); }); }); // 查询broker配置 ListConfigEntry entries admin.describeConfigs( List.of(new ConfigResource(ConfigResource.Type.BROKER, 0)) ).all().get().values().iterator().next().entries(); entries.forEach(e - System.out.println(e.name() e.value())); }可视化工具方面我用过几款简单说说区别。Kafka UIkafka-ui是开源里体验较好的能看topic、consumer lag、broker状态还有消息浏览功能。Offset Explorer原Kafka Tool比较传统图形化查看分区的offset和lag适合快速排查。实际定位问题时我更喜欢直接用AdminClient脚本配合命令行工具因为可视化工具在某些网络环境下会持续报连接异常反而干扰判断。排查“丢消息”时最核心的监控指标就是消费者的lag值。lag是消费者当前offset和log end offset之间的差值。如果lag持续增长说明消费跟不上生产如果lag突然清零但业务数据短了一截大概率是offset被提交跳过了数据。用Kafka UI构建端到端的lag监控是每个Kafka集群的必修功课。4. 消费端多线程、顺序性与延迟问题把消费端单独拎出来写一节是因为我见过太多团队在“消费端多线程”“顺序性”“消息延迟”这三个词上栽跟头。这三个问题表面上互不相干但经常在同一个项目里纠缠出现。4.1 多线程消费模式下顺序性该怎么保Kafka的顺序保证是分区级的——同一个分区内的消息消费者读取顺序与写入顺序一致。但如果你在消费端引入多线程顺序就很容易崩。先讲一个最简单的顺序消费方案单分区单线程。把partitions设置为1消费者用一个线程poll并且同步处理天然保序。但这种方式吞吐太低一般只用于对顺序极其敏感的场景。再讲实用方案分区级多线程。KafkaConsumer的poll方法本身是单线程的但你可以把拉取到的消息按分区hash到不同的处理线程ExecutorService[] executors new ExecutorService[partitionCount]; for (int i 0; i partitionCount; i) { executors[i] Executors.newSingleThreadExecutor(); } while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(1000)); records.partitions().forEach(tp - { ListConsumerRecordString, String list records.records(tp); executors[tp.partition() % partitionCount].submit(() - handle(list)); }); }核心思路是每个分区只对应一个固定线程消息进入分区时保持顺序线程处理时也按分区串行执行从而保证单分区顺序性。但要格外注意handle(list)里不能跨分区做共享可变状态。否则顺序虽然不崩业务结果还是会乱。有一个很容易被忽略的问题rebalance和线程池的配合。如果线程池里的任务还没处理完consumer实例已经触发rebalance已经poll出来但没处理完的消息会丢对当前消费组而言它会重新分配给其他消费者产生重复。我的建议是在处理完当前批次的poll记录后再调用commitSync并给线程池留出优雅停机的时间窗口。4.2 消息延迟高性能瓶颈到底在哪一端“Kafka消息延迟高”这个热搜词背后其实藏着两种完全不同的场景。第一种是生产写入延迟高。这种情况先看broker的IO压力用iostat看磁盘util和await。如果磁盘util持续超过80%说明日志写盘追不上生产速度。这时候需要做两件事一是检查topic分区数是否过少、生产者并发是否够二是检查是否执行了大量随机写比如多个broker共用一块磁盘上的其他业务。第二种是消费端处理延迟高。最典型的表现是lag不断上涨。先别急着加消费者线程应当先看消费逻辑本身有没有瓶颈。比如消费端调用了外部HTTP接口平均耗时2秒那再多的线程也可能不够。另一种隐蔽的情况是消费者max.poll.records设置得太大单次poll拉回来几万条消息业务处理时间超过了max.poll.interval.ms默认300秒消费者被判定为死亡并触发rebalance。rebalance期间所有分区暂停消费lag瞬间飙升。排查时我习惯从下往上打点先看broker端磁盘和网络再看客户端poll耗时最后看业务处理耗时。哪里出现拐点瓶颈就在哪里。还有一类延迟高与“读写最大值”相关。Kafka是高吞吐系统但单条消息大小直接影响吞吐上限。如果你发送的消息体动辄几十MB网络和磁盘的传输量是指数级增长。这类场景需要单独设置message.max.bytes、replica.fetch.max.bytes、fetch.max.bytes三个参数。调优之前想清楚一条消息几十MB本身就不适合走Kafka直接上对象存储吧。5. 典型案例与排查速查表学再多的原理不如看几个真实的翻车现场。这里我挑两个记忆深刻的案例再放一张排查速查表方便大家直接抄作业。5.1 InvalidReceiveException和各种网络异常实录有段时间线上集群持续报org.apache.kafka.common.network.InvalidReceiveException: Invalid receive from ...生产端消息大量发送超时。乍一看有人以为是客户端版本不兼容或者是broker的socket.request.max.bytes设置太小。后来排查发现问题出在客户端和broker之间的TCP连接被负载均衡设备断开。部分请求包在传输层被截断服务端收到的报文长度异常于是抛出InvalidReceiveException。更诡异的是我们只在特定时间段出现这个问题后来发现是负载均衡那里的空闲连接超时时间太短连接空闲超过阈值就被重置。这个案例说明Kafka丢消息和网络基础设施的关系极其紧密。遇到这类报错排查顺序是防火墙规则、负载均衡超时策略、客户端到broker的TCP缓冲区、Kafka协议版本兼容性。不要一上来就改Kafka参数。另一个让我记忆深刻的案例是跨机房双活场景。生产者A机房到消费者B机房专线抖动、丢包率上升。由于Kafka是长连接加上没有开启ack确认机制生产端以为发送成功实际broker收到的是不完整的数据包。排查后发现Kafka客户端默认的socket.send.buffer.bytes和socket.receive.buffer.bytes是百K级别高延迟专线下很容易触发拥塞。解决办法是把这两个参数调大同时启用重试机制。5.2 从面试题看Kafka丢消息的知识盲区我发现面试题往往能精准暴露一个人对Kafka理解深浅。这里整理几个高频问题也是在自查你踩过哪些坑。“为什么Kafka吞吐量高却还可能丢消息”这是一个好问题。Kafka高吞吐源自顺序写和零拷贝但吞吐高不等于可靠。吞吐优化和可靠性优化在某种程度上是对立的为了吞吐你可能关闭幂等、调低acks、增加批量大小为了可靠性你又要开启幂等、acksall、减少批量。两者之间需要业务权衡。“消费端多线程怎么保证消息顺序性”答案就是按分区把消息路由到固定线程单分区内串行消费。跨分区之间的顺序本来就不保证谁拿这个说需求就是外行。“消息重复消费怎么解决”核心是幂等。Kafka本身不能完全避免重复消费尤其是在rebalance和消费者网络抖动场景下。最简单的幂等方案是用业务唯一ID比如订单号做去重Redis的SETNX就能搞定大部分场景。“Kafka读写最大值与硬件有什么关系”Kafka的读写最大值主要受磁盘顺序写速度、网络带宽、内存影响的page cache命中率三者约束。NVMe固态硬盘的随机读写和顺序写都远超机械盘是部署Kafka的首选。内存越大page cache命中率越高读取消息时的IO越少。实际压测时可以按“磁盘顺序写带宽”和“网卡带宽”两个瓶颈来估算峰值吞吐。5.3 避坑经验总结我把这些年踩过、看过的坑汇总成一张速查表建议贴在你的运维文档里现象可能原因对策消息少了生产者acks0/1改成acksall消息少了topic副本数1设置replication.factor3消息少了broker断电没刷盘调整刷盘参数使用独立NVMe消息少了unclean.leader.election.enabletrue设为false消息少了消费端自动提交offset改为手动提交先处理再提交消息重复消费端无幂等处理业务幂等校验消息延迟高消费者rebalance检查max.poll.interval.ms与处理耗时消息延迟高磁盘IO瓶颈换SSD、增加分区InvalidReceiveException协议版本/网络设备截断检查版本匹配与负载均衡策略这里我必须强调速查表只能帮你快速定位真正做可靠性设计时一定要在发送、存储、消费三个环节同时加固。单独改一个参数永远不够。最后再分享一个小技巧每次发布Kafka相关代码或配置变更我都会先写一段“生产探针”往一个专门用于验证的topic里发送带序号的消息消费端验证序号连续性和延迟。发布完跑五分钟没问题才敢脱手。这套方法帮我挡住了至少三次线上事故你们也可以试试。
阅读完成 · 觉得有帮助?