ARTICLE DETAIL

资讯详情

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

Kafka消息不丢失:从原理到生产配置全解析

Kafka消息不丢失:从原理到生产配置全解析 上周帮一个朋友排查Kafka丢消息的问题现象特别典型生产端显示发送成功消费端却一直等不到数据。翻了一圈配置发现acks设成了1副本因子也是1一个broker宕机消息就悄悄没了大家还在群里争论到底哪一环出了问题。这个事让我意识到很多关于Kafka消息不丢失的讨论都停留在概念层真正到了线上配置对不上号的细节太多了。所以我想把这件事一次拆透先梳理一条消息从生产到消费到底要过哪些环节再给出生产环境可落地的配置组合和代码写法最后把高频面试题和一次真实事故复盘一起打包。正在学Kafka、准备自己搭集群的人以及被这类面试题卡过的同学读完应该能直接拿去用。1. 先别急着写代码搞清“消息不丢失”到底在防什么1.1 一条消息从生产到消费要过哪几道关很多人一上来就谈参数但参数只是在某个环节堵漏洞不先看清整条链路后面很容易出现“这里没丢、那里丢了”的尴尬。一条消息在Kafka里的完整路径大致要经过四关第一关生产者发送。Producer把消息放到自己的缓冲区里攒一批再通过网络发出去。第二关Broker接收。消息到达某个分区对应的Leader副本Leader把数据写入本地日志。第三关副本同步。Follower副本从Leader拉取这条消息拉取成功后消息才算真正进了ISR集合。第四关消费者消费。Consumer主动拉取消息执行完业务逻辑之后再提交消费位移。这四关里面任何一环出问题都可能造成“消息丢了”。我通常喜欢拿寄快递来类比Producer是寄件人Broker是快递中转仓库分区副本就是仓库里多留的几份底单Consumer是收件人位移提交则是收件人签收的动作。如果寄件人把包裹交给了快递员但没有要回执快递员路上把件弄丢了寄件人完全不知道——这就是acks1的感受。如果仓库只留了一份底单仓库着火底单没了那就真没了——这就是副本因子为1的处境。如果收件人还没验货就签收事后发现货物破损也没办法找快递公司索赔——这就是自动提交位移的问题。所以我们说Kafka消息不丢失从来不是某一个参数单独保证的而是生产端、Broker端、消费端共同配合出来的结果。1.2 三个环节、三类丢失责任边界要分清我经常在群里看到有人问“Kafka为什么会丢消息”接着吵成一团。其实只要把丢消息的场景按环节分个类责任边界就清楚了。环节常见丢消息原因对应关键配置Producer侧发送失败未重试缓冲区溢出直接丢请求acks1/0导致Leader返回成功但数据未同步acks、retries、delivery.timeout.ms、enable.idempotenceBroker侧副本因子为1ISR里的副本全部宕机Leader选举出非同步副本刷盘策略过于激进replication.factor、min.insync.replicas、unclean.leader.election.enableConsumer侧自动提交位移导致未处理就提交处理过程中进程崩溃Rebalance期间位移没提交enable.auto.commit、手动commitSync、ConsumerRebalanceListener这里有个很重要的认知Kafka本身的高可用设计是靠谱的但它是按默认配置跑的还是按“不丢消息”的标准配置跑的结果完全不同。我见过太多团队直接把默认配置丢到生产出了事就喊“Kafka丢数据”实际上Kafka是在按你的配置执行。默认的acks1、默认的自动提交本身就是为了性能和易用性牺牲了一部分可靠性。所以聊消息不丢失先别急着甩锅把三个环节的责任边界画出来再逐项配置才是正确的打开方式。2. 部署环节决定下限集群、副本与存储配置一次到位2.1 用KRaft模式把集群搭起来版本与服务端规划部署是消息不丢失的第一道分水岭。很多人的Kafka集群只有一台机器、一份数据那不叫集群叫单点。单点宕机意味着数据直接没后面聊再多参数都白搭。先谈版本。现在搭Kafka建议直接用KRaft模式也就是去掉ZooKeeper、用Raft协议管理元数据的架构。从Kafka 3.7开始KRaft模式已经很稳3.9之后官方逐渐把KRaft作为主推方向4.x已经默认走KRaft。你再去搜“ZooKeeper部署Kafka”这类老教程可以看原理但新环境不建议照搬。KRaft模式下节点分成Controller和Broker两种角色Controller负责元数据和分区Leader选举Broker负责实际的消息读写。小集群可以让一个进程同时跑Controller和Broker大集群建议分离。下载选择上直接去Apache Kafka官网下载二进制包就行解压即用。需要注意Java版本Kafka 3.x基本要求Java 11以上4.x要求Java 17以上环境变量配不对会直接起不来。Windows环境下Win11上想本地玩集群的话我更推荐用Docker Desktop配合WSL2或者直接在WSL2里跑二进制包避免Windows原生环境下的脚本兼容性和路径问题。下面给一个最简的docker-compose示例使用KRaft模式单Controller单Broker适合本地开发和练手services: kafka: image: apache/kafka:3.9.0 container_name: kafka ports: - 9092:9092 environment: KAFKA_NODE_ID: 1 KAFKA_PROCESS_ROLES: controller,broker KAFKA_CONTROLLER_QUORUM_VOTERS: 1localhost:9093 KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093 KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER KAFKA_CONTROLLER_LISTENER: CONTROLLER://:9093 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 volumes: - kafka_data:/var/lib/kafka/data - kafka_meta:/var/lib/kafka/meta volumes: kafka_data: kafka_meta:这里有个最容易踩坑的点advertised.listeners。它的作用是告诉客户端“你该用什么地址连我”。如果你在容器里只写了监听器地址没有写advertised客户端拿到的可能是容器内部IP从宿主机或另一台机器连不上。上面示例里我用的是localhost本地玩没问题跨机器部署就得改成对外的IP或者域名。如果你的环境启用了SSLlistener.security.protocol.map里要区分PLAINTEXT和SSL还要把keystore、truststore证书挂载进容器并把advertised的协议也改成对应的监听器。证书这一块很多新手漏配结果外部客户端连接时直接握手失败现象也是“连不上”、“拉不到数据”。2.2 副本因子、ISR与同步副本一切配置的起点集群搭起来之后接下来要做的不是急着写生产者代码而是把topic层面的数据冗余规则定好。这条规则由三个参数共同决定replication.factor分区副本数建议生产环境至少3。min.insync.replicas分区“最小同步副本数”意思是只有至少这么多个副本同步了消息才算提交成功建议设为2。acks生产者等待确认的条件不丢消息的标配是all。先解释一下ISR。ISR全称In-Sync Replicas同步副本集合。Leader副本是分区的写入入口Follower副本会不停从Leader拉取数据追上进度之后就在ISR里。如果有人落后太多比如网络故障、GC停顿Leader会把它从ISR里踢出去等它追上了再拉回来。这是一个动态过程不是配完就不变的。为什么说这三件套是一套组合拳replication.factor3意味着每个分区有3份数据任何1台Broker宕机另外2台还有完整数据。min.insync.replicas2则意味着一条消息要至少被2个副本都写成功了才向生产者返回成功。配合acksall生产者会等待所有ISR成员都确认写入才认为发送成功。这个组合下即使某个副本突然宕机消息仍然留在另外的同步副本里不会丢。代价是什么可用性下降。如果集群里只剩下1个可用副本而min.insync.replicas2生产者再发消息就会报NotEnoughReplicas异常生产失败。这是有意为之宁可写入失败也不让消息悄无声息地丢。很多人问“为什么我Kafka写入报错”先查一下是不是ISR里的副本不够。这是“不丢”换来的必然代价必须让业务方接受。还有一个关键参数unclean.leader.election.enable。如果ISR里的所有副本都挂了而一个之前落后很多、数据不完整的副本还活着要不要让它当Leader设为true的话集群能继续用但会丢掉ISR副本里那些它没有的数据这就是真正意义上的丢消息设为false的话这个分区直接不可读写等服务恢复。对“不丢消息”有硬性要求的场景必须设成false。宁可短暂不可用也不能让消费端读到残缺的数据。2.3 Docker部署最容易埋下的三个坑Docker已经把Kafka的安装门槛降得很低了我经常跟人开玩笑一条docker run就能起一个Kafka但“能起来”和“能用得稳”是两码事。部署阶段最容易埋下的隐患有三个全是血泪经验。第一个坑没挂载数据卷。如果你只是写了image、ports靠容器内部目录存数据一旦容器被删数据全部归零。Kafka在KRaft模式下的数据分为两部分业务消息数据在/var/lib/kafka/data元数据在/var/lib/kafka/meta。这两个目录都要挂载到宿主机持久化目录否则容器升级、迁移、重建的时候等着你的就是一个全新的空集群。这个坑在Windows Docker Desktop上更隐蔽很多人把目录挂载到Windows文件系统后发现Kafka因为权限问题起不来最终还是要落到WSL2的文件系统里。第二个坑advertised.listeners配置不对。前面提过容器里监听地址是0.0.0.0:9092但对外广告的地址必须是你客户端真正能访问的地址。如果你在Docker端口映射里把9092映射到宿主机9092advertised写成PLAINTEXT://localhost:9092本机连没问题但跨机器就不行得写成宿主机IP。很多人在服务器上部署客户端在另一台机器连不上查来查去最后发现就是advertised没改。第三个坑内存不受控。Kafka是JVM应用不设置堆大小就会用默认值容器memory limit给得小JVM直接OOMKilled容器反复重启如果不设limit宿主机内存又被拖垮。推荐在容器环境变量里设置KAFKA_HEAP_OPTS-Xms2g -Xmx2g容器memory至少给3到4G。不要以为小集群不需要内存Kafka的页缓存非常吃内存只是平时不容易被看到而已。3. 生产者与消费者配置把“不丢”落到代码上3.1 生产者端acksall、重试与幂等是标配部署和topic配置属于“地基”生产者和消费者的客户端配置则是“梁柱”。先看生产端。acks有三个可选值理解它们的区别就理解了Kafka丢消息最常见的一个源头acks0发出去就算成功不等待任何确认。吞吐最高但消息在网络上丢了、Broker挂了生产者完全无感知。我只会建议用在日志上报这种允许丢弃的场景。acks1Leader副本写入本地就返回成功。大多数默认配置就是这个值性能和数据安全平衡但Leader在返回之后、Follower还没来得及同步时宕机这条消息就没了。acksall生产者会等待分区ISR里所有副本都写入成功后才收到成功响应。这是“不丢消息”必须的起点。很多人以为设了acksall就万事大吉其实还不够。重试机制也要配合好。生产者在发送失败时会重试但如果retries配得太小或者delivery.timeout.ms太短重试还没成功就超时放弃了。我在生产环境一般这样配Properties props new Properties(); props.put(bootstrap.servers, kafka1:9092,kafka2:9092,kafka3:9092); props.put(key.serializer, org.apache.kafka.common.serialization.StringSerializer); props.put(value.serializer, org.apache.kafka.common.serialization.StringSerializer); props.put(acks, all); props.put(enable.idempotence, true); props.put(retries, Integer.MAX_VALUE); props.put(delivery.timeout.ms, 120000); props.put(max.in.flight.requests.per.connection, 5); props.put(batch.size, 16384); props.put(linger.ms, 20); props.put(buffer.memory, 33554432);重点解释两个容易被忽略的选项。enable.idempotencetrue是幂等生产者它会为每个消息带上生产者和序列号Broker端用来去重这样即使网络超时重发也不会造成重复消息。Kafka 3.0之后当你设置acksall时幂等默认开启但显式写出来更清楚。max.in.flight.requests.per.connection控制在单个连接上最多有几个未确认请求5是幂等开启时的默认值可以保证消息顺序的同时兼顾吞吐。batch.size和linger.ms对“丢”没有直接影响但会影响性能。有人为了让消息快点发出去把linger.ms设成0结果每个消息都单独一个请求吞吐很难看反过来为了攒批把linger.ms设得过大消息延迟又上去了。我在生产环境一般用20ms左右具体根据自己的业务响应要求调。3.2 消费者端手动提交位移把主动权握在自己手里生产端把消息送进Kafka只是第一半另一半在消费端。消费者丢消息最常见的元凶是自动提交位移。Kafka默认的enable.auto.committrue意思是Consumer在poll返回一批消息之后会自动把当前位移提交上去默认5秒一次。问题来了如果自动提交发生在你处理完业务之前进程突然崩溃重启后Kafka以为你已经消费完了就从已提交的位移继续往后拉中间那批没处理的消息就丢了。所以“不丢消息”的消费端一定要这样做Properties props new Properties(); props.put(bootstrap.servers, kafka1:9092,kafka2:9092,kafka3:9092); props.put(group.id, order-service); props.put(key.deserializer, org.apache.kafka.common.serialization.StringDeserializer); props.put(value.deserializer, org.apache.kafka.common.serialization.StringDeserializer); props.put(enable.auto.commit, false); props.put(max.poll.records, 500); KafkaConsumerString, String consumer new KafkaConsumer(props); consumer.subscribe(Collections.singletonList(orders)); while (running) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecordString, String record : records) { // 先处理业务入库、调用外部接口等等 process(record); } // 业务处理完再手动提交位移 consumer.commitSync(); }这里面的顺序很关键先处理再提交。这样即使提交前进程挂了重启后会从旧位移重新拉取消息会重复消费但不会丢。Kafka默认的语义就是at-least-once也就是“至少一次”不保证不重复。所以业务逻辑要做好幂等这是另一码事。commitSync和commitAsync怎么选commitSync会阻塞并自动重试适合在一个批次处理完之后调用commitAsync不阻塞提交失败时不会自动重试但可以通过回调记录失败。我常用的组合是每次批次处理完用commitAsync异步提交提高吞吐在进程关闭前最后再用一次commitSync把未提交的位移同步提交掉确保优雅退出时不丢位移。另外要留意Rebalance。当一个消费者进程挂了或者新消费者加入分区会重新分配。如果ReBalance发生的那一刻正在处理的消息还没提交位移这批消息就会被重新分配给其他消费者重新消费。这是Kafka重复消费的最常见场景之一也是at-least-once语义的自然结果。如果业务对乱序和重复很敏感可以做两件事一是监听ConsumerRebalanceListener在revoke之前提交位移二是让处理逻辑尽量幂等比如按业务主键去重。3.3 关于“生产消费命令启动一次会一直运行吗”的澄清很多刚开始用Kafka的人会问一个很基础的问题kafka-console-producer和kafka-console-consumer启动一次会一直运行吗答案是它们都是常驻进程。生产者在启动后会等待你在终端输入消息不会自己退出消费者启动后会持续监听新消息并打印到终端也不会自己退出。想停止就按CtrlC。这个设计本身是为了模拟真实客户端的行为。你生产环境里的Java服务也是长期运行的不可能发一条消息就退出。控制台工具只是把客户端做成了命令行版方便测试和排查。如果你只是想快速验证topic里有没有数据又不想被持续刷屏可以加max-messages参数限制一次最多拉取多少条。我平时排查问题最常用的命令是# 查看topic列表 kafka-topics.sh --bootstrap-server localhost:9092 --list # 查看某topic的分区、副本、ISR等详情 kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic orders # 从头开始消费这个topic最多拉20条 kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic orders --from-beginning --max-messages 20注意那些-run-class形式的命令比如kafka.tools.GetOffsetShell本质也是一次性查询工具执行完会返回结果不是常驻的。“常驻”和“执行一次”在这个领域里不是矛盾关系而是工具分工不同。4. 从原理到监控把“丢没丢”变成可观测4.1 消息发成功了却查不到先从这几个地方找“生产端明明返回成功了消费端却查不到数据”这是我被问到最多的场景之一。多数时候消息没丢只是坐标没找对。第一种情况消费者组没有历史位移。一个新group去消费一个早已存在的topic默认auto.offset.reset是latest它只会消费启动之后的新消息之前的消息看起来就像“丢了”。你只要从头消费一次或者把auto.offset.reset设为earliest再重置位移旧数据就能看到了。第二种情况topic不存在或者分区Leader不在线。可以用kafka-topics.sh --describe --topic xxx查看分区和ISR状态如果提示leader为-1说明Leader选举有问题消息暂时不可读。第三种情况消费组的位移已经跑到了最后。想看历史数据可以直接用console consumer指定--from-beginning拉一遍。但要提醒一句生产环境的topic数据量一般都很大不要直接从头拉会冲击Broker网络务必加上--max-messages限制条数。如果只想看某个分区当前的消息位移可以用kafka-run-class.sh kafka.tools.GetOffsetShell --broker-list localhost:9092 --topic orders --time -1这个命令会返回每个分区最新的offset。再用kafka-consumer-groups.sh --describe --group order-service查看消费组当前的位移和lag两条命令一对比就能判断消息到底有没有被消费、消费到了哪个位置。排查此类问题的思路就是先确认数据在不在Broker上再确认消费者有没有把位移推过去这两步走完九成问题都能定位。4.2 消息延迟高和OOM为什么它们和“不丢”强相关消息延迟高本身不等于消息丢失但延迟长期堆积会引发一系列“看起来像丢消息”的连锁反应。最常见的情况是消费端处理不过来poll出来的消息积压处理耗时太长超过了max.poll.interval.ms消费者被判定为“死亡”触发Rebalance。Rabalance期间分区被重新分配负责处理旧分区的消费者还没提交位移消息就会被其他消费者重新消费一遍。这个过程中如果某个消费者进程直接崩溃那些未提交的消息就相当于在业务层“丢”了。线上排查延迟高的问题我习惯按下面这个顺序走先看Lag。kafka-consumer-groups.sh --describe --group xxx --bootstrap-server localhost:9092看LAG列是否持续增大。Lag大说明生产速度大于消费速度问题在消费端。再看磁盘IO。Kafka写消息依赖磁盘如果Broker所在机器的iowait高分区写入慢生产端会跟着超时、重试整体影响链路。再看GC。Kafka是JVM应用Full GC期间会stop-the-worldBroker暂停处理请求。GC停顿时间过长会导致Follower副本同步跟不上被踢出ISR甚至触发Leader切换。用jstat看GC频率和耗时如果老年代持续增长要么堆太小要么有内存泄漏。最后看网络和CPU。网络带宽打满、CPU跑满都可能让请求排队表现为延迟升高。OOM和消息不丢失的关系更直接。Kafka Broker一旦OOM崩溃该节点上所有分区的Leader都会重新选举如果此时ISR里副本数不够部分分区可能不可写如果崩溃前数据还没有同步到其他副本数据就直接没了。我见过不少Broker OOM的案例根子大多是堆内存设置不合理分区数却越加越多每个分区都占着对应的索引、缓存和请求队列堆逐渐被撑爆。建议根据分区数量和角色职责设置KAFKA_HEAP_OPTS一般小集群4G起步大集群8G甚至更高同时打开GC日志方便事后定位。KAFKA_HEAP_OPTS-Xms4g -Xmx4g -XX:UseG1GC -XX:MaxGCPauseMillis204.3 监控指标与ELK/OTel集成聊了这么多配置最怕的就是配完之后不知道系统实际状态。监控的作用是“在丢消息的第一时间发出告警”把隐性风险变成显性指标。我每次搭Kafka集群下面这几个JMX指标是必须盯的指标含义重点关注UnderReplicatedPartitions副本不同步的分区数长期大于0说明有副本跟不上丢消息风险高OfflinePartitionsLeader不在线的分区数大于0说明有分区不可读写IsrShrinksPerSec / IsrExpandsPerSecISR收缩/扩张频率频繁收缩说明Broker不稳定RequestHandlerAvgIdlePercent请求处理线程空闲率低于30%说明Broker繁忙BytesInPerSec / BytesOutPerSec入站/出站吞吐观察流量趋势辅助容量规划落地方式上我推荐用Prometheus加kafka_exporter采集指标再配合Grafana出面板。kafka_exporter会把JMX指标转成Prometheus格式告警规则可以写成UnderReplicatedPartitions大于0持续5分钟就报警同时配合consumer group的lag指标能直观看到消费堆积。至于ELK则是把Kafka的server.log、controller.log用Filebeat采集经过Logstash解析后写入Elasticsearch最后在Kibana里做日志检索。日志和指标配合使用排查问题会快很多。如果是微服务链路比较多的场景还可以考虑用OpenTelemetry它既可以采集Kafka的JMX指标上报到Prometheus后端也可以通过OTel的Kafka receiver直接消费Kafka里的指标和链路数据还能在客户端SDK里自动埋点把生产消费的延迟、错误率接入统一的可观测平台。运维监控上我建议宁可多配几个告警也不能让“副本不同步”这种信号在角落里躺几个小时。5. 面试考点与一次真实的事故复盘5.1 三种投递语义至少一次、最多一次、恰好一次Kafka面试题里“消息不丢失”这个话题几乎必考。要回答到位得先搞清楚投递语义的分类。at-most-once最多一次。允许消息丢失不允许重复。做法是消费者先提交位移、再处理消息处理失败就不管了重启后直接从已提交的位移往后走。适合那些丢了可以重算、重复会导致严重问题的场景。at-least-once至少一次。不允许消息丢失允许重复。做法是先处理消息、再提交位移处理失败后重启会重新消费所以不丢但可能重复。Kafka默认就是这种语义。exactly-once恰好一次。不丢也不重复。需要幂等生产者和事务机制配合成本高通常只在Kafka Streams这类流处理场景里用到。面试里如果问“Kafka到底会不会丢消息”正确的回答姿态是Kafka在正常配置下不会丢但前提是生产端、Broker端、消费端都按规范配置。它默认提供的是at-least-once语义既不保证“一定不丢”也不保证“一定不重复”。想不重复要么开幂等要么做业务幂等。5.2 一个真实的消息丢失事故复盘前年帮一个团队复盘过一次线上事故典型的配置不规范导致消息丢失拿出来讲讲希望你不要重走这条路。背景是三个节点的Kafka集群某个核心topic创建得比较早一直用的是默认配置replication.factor1min.insync.replicas1。某天其中一台Broker的磁盘损坏整个节点下线。由于topic只有一个分区副本正好落在这台机器上分区直接不可用。等运维把节点恢复后重新拉起发现该分区的数据已经无法完整读取丢了一批消息。更麻烦的是生产端用的是acks1Leader写入成功就返回客户端日志里看不到任何失败。排查过程其实不复杂。先看告警发现UnderReplicatedPartitions在事故前就已经是1持续了好几天说明这台Broker早就同步异常了但没人处理。再看Broker日志分区报错类似NoReplicaOnlineException一个在线副本都没有。最后查生产端配置acks1。三个线索放在一起结论很清楚副本因子为1导致数据没有冗余acks1又让生产者完全不知道数据没有被其他副本同步。修复措施分三步走。第一步把该topic的副本因子调整到3做分区副本重分配让数据分布到多个节点。第二步统一生产端配置为acksall并设置min.insync.replicas2这样以后即使某个副本挂了写入也会报错而不是静默成功。第三步在运维层面把“单副本topic巡检”做成常态化凡是新建topic强制要求replication.factor不低于3minISR不低于2从源头杜绝类似问题。这个事故给我最大的触动是Kafka的“高可靠”是配出来的不是装出来就自带光环。默认配置下它可能只是一个性能不错的消息管道距离“不丢消息”还差得很远。5.3 面试怎么答“Kafka如何保证消息不丢失”面试官问这个问题考察的往往不是你会背几个参数而是你有没有完整的链路意识。我建议按“三段式”回答逻辑清晰也不容易漏点。第一段说生产端设置acksall让Leader等待ISR中所有副本确认写入再返回成功开启幂等生产者enable.idempotencetrue防止重试导致的重复消息合理设置retries和delivery.timeout.ms给重试留足时间。第二段说Broker端topic副本因子至少3min.insync.replicas至少2保证数据有多副本冗余unclean.leader.election.enable设为false绝不允许非同步副本参与Leader选举避免丢数据换取可用性。第三段说消费端关闭自动提交位移enable.auto.commitfalse业务处理完成后再手动提交实现at-least-once语义配合幂等消费把重复问题也在业务侧消解掉。最后可以补一句Kafka本身通过ISR、HW高水位、LEO机制来定义“消息已提交”只有被ISR中多数副本写入后的消息才对消费者可见。在配置正确的前提下Kafka是能够做到消息不丢失的。这句话一说面试官就知道你是真懂原理不是背答案。6. 常见问题速查表与最终自检清单6.1 问题排查实录速查表排查Kafka问题最忌讳瞎猜。下面这个表格是我平时排除故障时常用的速查思路按“现象-原因-排查-处理”四列整理基本覆盖了消息不丢失相关的绝大多数场景。现象可能原因排查思路快速处理生产端显示成功消费端看不到消费者组从latest开始消费老消息没被消费查看group lag确认offset位置用--from-beginning或重置group offset到earliest消费者重启后重复消费位移提交晚于消息处理或关闭前没提交检查auto.commit配置和commitSync调用时机改手动提交关闭前最后同步提交一次生产端报NotEnoughReplicasminISR大于当前可用同步副本数查看分区ISR状态检查Broker存活恢复副本同步或评估风险后临时降低minISR消息延迟持续升高消费者处理慢、poll超时、磁盘IO高看lag、iostat、jstat、线程栈优化消费逻辑增加消费者实例修复Broker IO瓶颈Broker频繁OOMKilled容器内存限制过小、JVM堆溢出看容器events和GC日志调整KAFKA_HEAP_OPTS给容器多留内存容器重建后数据全没了数据卷未挂载或挂载目录错误查看docker inspect的Mounts挂载/var/lib/kafka/data和/var/lib/kafka/metaUnderReplicatedPartitions长期大于0某个副本长时间同步不上查看Broker日志、磁盘、网络踢出异常Broker恢复副本同步外部客户端连接不上advertised.listeners配置错误、SSL证书问题检查listener配置和证书有效期修正advertised地址重新生成证书6.2 部署与配置自检清单最后我把“Kafka消息不丢失”这件事浓缩成一份可以直接打勾的自检清单。每次新部署一个集群或者新接一个业务我都会先过一遍这份清单Broker端副本因子不小于3min.insync.replicas不小于2unclean.leader.election.enablefalseKRaft元数据目录和数据目录都已持久化容器内存和JVM堆大小匹配预留足够PageCache。生产者端acksallenable.idempotencetrueretries设置合理delivery.timeout.ms大于单次重试上限batch.size和linger.ms不拖垮延迟buffer.memory满足高峰期积压。消费者端enable.auto.commitfalse先处理业务再提交位移关闭前最后一次commitSync监听ConsumerRebalanceListener在分区被回收前提交位移消费逻辑具备幂等能力。监控告警UnderReplicatedPartitions、OfflinePartitions、ISR收缩、consumer lag、Broker GC五类指标都有采集和告警日志已接入ELK或类似平台可检索链路观测已通过OpenTelemetry接入统一可观测平台。我现在每次搭Kafka集群第一步不会急着写业务代码而是在脑子里把整条消息链路过一遍从生产、Broker、消费到监控逐项对照这份清单打勾。打勾的过程不复杂但它能避免绝大多数“消息丢了”的烂事。最后再分享一个小习惯凡是新建topic我都会先用kafka-topics.sh --describe看一眼副本数和ISR状态再让业务接入。这个习惯帮我躲过好几次线上事故你也不妨试试。
返回列表