
搞Kafka的人几乎都躲不过这四大灵魂拷问消息到底会不会丢、顺序能不能保证、消费会不会重复、消费者组为什么又又又重平衡了。这四个问题不光是高频面试题更是生产环境里的“四大天王”——任何一个爆雷轻则数据对不上重则核心链路直接脑卒中。网上讲这些问题的文章不少但大多是把几个参数列一下完事没把背后的来龙去脉讲透容易看完就忘出了事还是不会查。我打算换个思路把这五个问题消息不丢失、顺序传输、不重复消费加上重平衡Rebalance串在一条完整的消息链路上拆开讲一条消息从生产者写进去到Broker存下来再到消费者拉走每一个环节都可能变形。搞懂了这条链路里谁可能丢、谁可能乱、谁可能重复你再看那些参数配置就是水到渠成的事根本不用死记。标题里的reblanace我猜是Rebalance的笔误不重要盘它就完了。这篇文章适合刚入门想建立完整认知的兄弟也适合被线上问题折腾过、想系统梳理一遍的老手。1. 先把链路画出来一条消息从生产到消费到底经过哪些关卡要聊清这四个问题得先把Kafka的基本地盘捋清楚。很多人对Kafka的理解停留在“一个高性能消息队列”这个层面但真正到了排查问题的时候连“谁是Leader、谁是Follower”都说不清那后面所有的原理都是空中楼阁。1.1 为什么Kafka能扛住高吞吐日志追加模型与分区并行Kafka性能好的核心在于它的存储模型极度简单。一个Topic被拆成若干个Partition每个Partition在物理上就是一个追加写的日志文件消息只往后写不改前面的数据。这个设计让Kafka能充分利用顺序磁盘IO随机读写在这个模型里几乎不存在。更关键的是Partition之间天然并行一个Topic有12个分区理论上就能让12个消费者线程各自消费一个分区互不干扰。这里有个很多人混淆的概念Partition数量决定的是并行度上限而不是存储容量。有人一看到消息积压就疯狂加分区其实加分区只是让更多消费者能参与进来如果消费者的处理逻辑本身有瓶颈分区再多也是白搭。Producer发送消息时会通过分区器选择目标分区。默认情况下如果消息带了Key就对Key做哈希取模保证相同Key的消息永远进同一个分区不带Key就轮询或者按照粘性分区策略分配。1.2 消息在Kafka中的完整生命周期我把一条消息从生到死的完整路径写出来后面所有问题都围绕这条路径展开Producer调用send()消息进入生产者客户端的缓冲区accumulator攒批后由Sender线程发往Broker。Broker收到消息后先落到对应分区Leader副本的本地日志然后向ISRIn-Sync Replicas同步副本集合中的Follower发起复制。Producer收到ACK表示消息写入成功具体等几个副本确认看acks配置。Consumer通过poll()主动拉取消息处理完业务逻辑后提交Offset到内部Topic__consumer_offsets。下次poll()时从提交的Offset之后继续拉取。注意Consumer是“主动拉”模式不是Broker“推”给Consumer的。这个设计意味着消费快慢完全由Consumer自己控制但也导致了很多问题——比如你处理不过来到底算谁的锅后面讲重平衡的时候会重点展开。1.3 理解两个关键角色ISR与CoordinatorISR是Kafka保证数据不丢的基石。一个分区的Leader副本维护着一个“同步副本集合”里面是所有跟得上Leader进度的Follower。消息写入Leader后只有ISR里的副本都复制成功了才算真正安全。Coordinator协调者则是消费者组的大脑。每个消费者组都会选出一个Broker作为Coordinator负责组成员管理、Offset存储和重平衡触发。所有消费者加入、退出、心跳都要跟Coordinator打交道。很多人排查重平衡问题第一步就是去看协调者在哪台机器上因为Coordinator所在节点的负载直接影响整个组的稳定性。这里可以做个比喻ISR像“存款保险”保证钱消息不会因为单一银行副本倒闭而蒸发Coordinator像“小区物业”管着谁住哪套房分区分配谁搬走了成员退出都要重新安排。把链路和角色搞清楚接下来就可以一个一个地拆四大问题了。2. 消息不丢失三个环节里最容易出错的不是Broker而是生产者和消费者Kafka丢消息很多人第一反应是“Broker挂了吧”但实际上生产环境中绝大多数丢消息问题出在两端——要么是生产者以为发出去了要么是消费者以为处理完了。Broker本身在默认配置下反而相当皮实。2.1 生产者端acks和retries配置不当消息会“静默”丢失先看Producer写入Broker的确认机制这里有一个最关键的参数acks。acks0Producer发出去就不管了不等待任何确认。最极端的情况是消息根本没到Broker你在日志里也看不到任何报错。这种配置下丢消息没有任何提示适合丢几条无所谓、追求极致的吞吐的场景比如日志采样、统计类数据。acks1Leader写入本地日志就算成功。这里有个坑如果Leader刚写完还没来得及把消息复制给Follower就宕机了并且新的Leader没有这条消息那这条消息就永久丢失了。从“可靠”角度看acks1只能算半可靠。acksall或-1Leader要等ISR里所有副本都确认写入后才返回ACK。这是最可靠的级别但也是延迟最高的级别吞吐会下降一些。生产环境如果对数据有明确要求我建议无条件上acksall。很多团队为了追求吞吐把acks设成1结果碰上Leader宕机副本切换消息悄无声息地没了数据对账的时候才哭。再说retries。它决定了消息发送失败后的重试次数。这里要强调一个细节重试不是简简单单“失败就重发”如果开了重试消息可能会在客户端内存里滞留如果你同时设置了delivery.timeout.ms两者是配合工作的。实际中我见过很多“丢消息”案例本质是发送过程中抛出了超时异常代码里catch住吞掉了然后没有补偿机制消息就没了。所以生产者端真正要防的不是Kafka丢而是应用代码把发送失败当成“可忽略异常”。2.2 Broker端什么时候会丢数据什么时候不会Broker端丢数据主要集中在两个场景。第一个场景是Leader副本宕机而Follower的复制落后了。只要Follower在ISR里就说明它跟Leader的差距在容忍范围内新Leader顶上后数据是完整的。但是——如果连ISR里的副本都全挂了或者因为各种原因ISR被缩到只剩一个副本那就存在数据窗口。这时候如果配置了unclean.leader.election.enabletrue意味着允许ISR以外的“落后副本”参与Leader选举这个落后副本可能没有最新数据选它当Leader就会丢消息。我强烈建议生产环境把这个参数保持false宁肯短暂不可用也不要丢数据。第二个场景是磁盘问题。消息写入PageCache后最终要刷到磁盘才安全。实际上Kafka依赖操作系统刷盘机制同时自身在配置里有flush相关的参数。正常情况下机器正常重启PageCache里的数据会刷到磁盘不会丢但如果是断电、宕机、云主机被强制停止就可能丢一部分PageCache数据。要彻底规避就得调log.flush.interval.messages和log.flush.interval.ms但这会影响性能。我的取舍原则是核心交易数据可以降低刷盘间隔日志类数据让它慢慢刷。2.3 消费端最容易忽略的“丢消息”其实是“没消费完就提交”消费端丢消息的场景更隐蔽。Kafka里有一个核心参数enable.auto.commit默认是true也就是说Consumer每5秒自动提交一次Offset。你想想这个模型的隐患如果你拉取了一批消息正在处理第1条时自动提交了这批消息的Offset——也就是说这个批次里即使有消息还没被处理Offset也已经标记成“已消费”了。此时消费者宕机重启后从已提交的Offset继续消费中间那些没处理完的消息就永远跳过了。这就是经典的“at-least-once vs at-most-once”问题。你以为你在用at-least-once至少一次但默认配置下你的消费语义实际上是at-most-once至多一次——至少会丢一部分消息。而开启手动提交enable.auto.commitfalse之后又面临另一个难题到底什么时候提交Offset处理完就提交万一提交成功但后续业务逻辑出错呢这里先不展开丢消息和重复消费往往是一体两面你在解决一个问题的同时会制造另一个问题。2.4 生产环境推荐配置一张表说清环节参数推荐值原因Produceracksall等待ISR内所有副本确认最大程度避免单点故障丢数据Producerretries大于0建议3-5网络抖动、Leader切换时自动重试Producerenable.idempotencetrue与acksall配套保证发送幂等性下一节细讲Producerdelivery.timeout.ms120000给重试留足时间窗口防止超时截断重试Brokerunclean.leader.election.enablefalse不允许非同步副本参与Leader选举宁可不可用也不丢数据Brokermin.insync.replicas大于1建议2保证至少存在一个Follower同步acksall才有意义Consumerenable.auto.commitfalse手动提交把控制权握在自己手里Consumerauto.offset.resetearliest无提交Offset时可重放配合幂等消费可以安全重读这里有一个经常被忽视的组合关系acksall和min.insync.replicas2必须同时配置。如果min.insync.replicas是1那么Leader自己就是ISR里的那个“1”acksall实际上退化成acks1等半天等于白等。如果ISR的副本数小于min.insync.replicasBroker会直接拒绝写入报出NotEnoughReplicasException——这也是判断副本健康度的一个信号。3. 顺序传输Kafka只保证分区内有序但“业务有序”比“技术有序”复杂得多顺序问题是我见过引发最多“哲学辩论”的话题。Kafka官方文档说得很清楚Kafka只保证单个分区内消息有序不保证一个Topic内全局有序。但很多人实际遇到的需求是“业务上的顺序”这就比单纯的技术顺序要复杂一个数量级。3.1 分区内有序的底层保证为什么分区内能保持顺序因为一个Partition在物理上是一个追加日志Producer往同一个分区发消息Broker只是顺序写入。对应地Consumer在拉取一个分区时也是顺序返回的Offset小的先返回。这里有一个关键点同一条业务链路上的消息如何保证进入同一个分区靠的就是消息Key。Producer在发送时指定Key比如订单ID相同Key的消息经过哈希后永远进入同一个分区。分区器默认用的是murmur2哈希不是简单的取模但效果一样——一个Key对应一个确定的分区。我遇到过一个朋友他把业务订单号作为Key但订单下的子任务分散到多条线程并发发送认为“反正Key一样会进同一分区”。这里就有一个隐藏问题虽然消息进了同一个分区但多线程并发发送时在客户端内部的顺序并不能保证和业务发生顺序一致。Kafka保证的是“你发送的顺序”在分区内的顺序而不是“业务发生的顺序”在分区内的顺序。这个概念绕但很关键。3.2 生产者重试会怎样打乱顺序即使你单线程发送重试机制也可能把顺序打乱。Kafka客户端引入了max.in.flight.requests.per.connection这个参数它限制了一个连接上最多有几个未确认的请求在飞行。这个参数和顺序的关系是如果设置大于5且开启重试第一个批次的请求失败后重试而第二个批次已经发送出去并被Broker接收那么两个批次到达Broker的顺序就反了。反过来如果把这个参数设成1那么同一连接上只有第一个请求完成了才会发第二个自然就保证了顺序代价是吞吐下降。如果你设置了enable.idempotencetrue开启幂等情况又不一样了幂等生产者内部有序列号机制即使max.in.flight.requests.per.connection大于1Broker也会根据序列号重新排列批次保证乱序消息不会真正写入。这就是为什么开启幂等后你可以把in-flight调高而不担心乱序的原因。3.3 既然Kafka不保证全局有序业务上怎么办很多业务需求动辄要求“全局有序”比如“所有用户操作都要按顺序处理”。这种需求在Kafka里用单分区方案可以实现但吞吐会限制在单分区的极限。更合理的方案是按业务维度拆分比如电商订单场景用户对同一个订单的“创建、支付、发货、完成”这些事件必须有序但不同订单之间互不干扰。那么消息Key就用订单ID订单ID相同的事件进同一分区天然有序不同订单分散到多个分区并行处理。这是我见过使用最广泛的模式。还得考虑消费者端的顺序。消费者拉到一个分区的消息后如果业务处理是多线程的比如用线程池并发处理顺序就又被打破了——哪怕你从分区里按顺序拉出来处理顺序可能乱套。所以对顺序有严格要求的消费逻辑最好保持单线程消费或使用有序的线程模型。这一点在Kafka Streams里有专门的设置max.task.idle.ms、线程模型在生产环境中真要追求顺序往往不是Kafka的问题而是消费者自己把自己搞乱了。4. 不重复消费为什么Kafka无法从架构上彻底消灭重复只能靠业务侧兜底很多人一开始都天真地以为Kafka作为成熟的消息中间件应该保证不丢不重。现实是Kafka能保证不丢在合理配置下但无法保证不重。原因要从它“至少一次”的交付语义说起。4.1 “不重复消费”的三个重复源头源头一生产者重试导致Broker重复写入。Producer发消息超时后重试如果第一次实际已经写入成功第二次重试就会造成Broker里出现两条一模一样的消息。这是最常见的第一类重复。源头二消费者处理成功但提交Offset失败。消费者拉取一批消息处理完业务逻辑正准备提交Offset消费者进程崩溃。重启后从旧Offset重新拉取这批消息会被重新消费一次。这是第二类重复。源头三消费者处理成功提交Offset成功但业务逻辑的“成功”是异步的比如写了数据库但事务还没提交又发了个通知之类的。这类重复更隐蔽往往发生在流程复杂的业务中很难定位。4.2 幂等生产者到底解决了什么Kafka在0.11版本引入了幂等生产者enable.idempotencetrue它解决的是第一类重复。原理是每个Producer在初始化时会被分配一个唯一的PIDProducer ID发送每条消息时带上一个单调递增的Sequence Number。Broker端为每个PID维护这个序列号如果收到重复的序列号直接丢弃。但注意两个边界。第一幂等只能保证单PID、单会话内的不重复一旦Producer重启PID变了之前的序列号作废重复写入问题就回来了。第二它只保证“写入不重复”不可能感知消费者是否重复消费。所以幂等生产者只是链路中的一环不是终点。如果把场景升级到跨Producer会话甚至跨实例比如事务性写入Kafka提供了事务APItransactional.id Producer事务它可以把多个分区的写入变成一个原子操作并配合read_committed隔离级别让消费端只读取已提交的消息。这套机制再配合消费者的幂等处理才能接近端到端的“恰好一次”。4.3 消费端幂等的几种实用套路既然架构上无法彻底避免重复业务侧的幂等就是最后的防线。我做过不少消费程序总结下来有三种比较实用的方案第一种是唯一键去重。消费消息时把消息ID或业务唯一键写入数据库唯一索引如果插入冲突说明已经处理过跳过。这个方案最简单但对写入场景比较友好更新类场景要配合版本号。第二种是状态机校验。针对订单这类状态流转清晰的业务每条消息都携带订单当前状态和目标状态消费时先查数据库判断当前状态是否等于消息里的前置状态不匹配就说明这条消息过期或者重复直接丢弃。这种方法把自己从“判断是否重复”中解放出来只判断是否“合法”。第三种是Redis与DB配合的分布式锁或幂等表。处理前先抢锁或记录处理标记防止并发重复消费。这个方案适合处理链路比较长、不适合用唯一索引的场景。我在生产实践中见过一个比较经典的组合拳Protobuf消息里塞了一个全局的唯一请求ID消费端用这个ID作为数据库主键同时把enable.auto.commit设为false只有在业务逻辑处理成功后才手动提交Offset。这样即使重复消费数据库主键冲突直接拦截不会产生脏数据。看上去简单但真的能扛住绝大多数重复场景。5. 重平衡Rebalance一次让整个消费者组停摆的“重新分配”触发原因与规避手段Reabalance是Kafka里最让人血压升高的机制之一。简单说当一个消费者组里的成员发生变化时Kafka会把分配给所有成员的分区重新洗牌这个过程叫Rebalance。在Rebalance期间整个消费者组都无法消费消息相当于集体暂停——你数据没丢但业务停了。5.1 触发Rebalance的四种主因第一种消费者主动加入或退出。比如新实例启动加入同一个group.id或者某个消费者进程优雅退出调用close方法。这类重平衡通常是计划内的影响可控。第二种消费者心跳超时。Consumer需要周期性地发送心跳给Coordinator由heartbeat.interval.ms控制频率如果Coordinator在session.timeout.ms时间内没收到心跳就会判定消费者死亡踢出组并触发Rebalance。造成心跳超时的原因很多消费者卡在长任务里、GC停顿、网络分区等。第三种消费处理超时。Kafka在0.10.1版本引入了max.poll.interval.ms参数默认300秒。意思是如果Consumer两次调用poll()的间隔超过5分钟Coordinator就认为这个消费者“拉肚子”了处理能力跟不上主动把它移出组触发Rebalance。这个机制的存在是因为光靠心跳不足以判断消费者是否真的活着——你心跳正常可能只是卡在一个重型业务逻辑里。第四种订阅的Topic分区数变化或者订阅的Topic列表发生变化。比如给Topic加了分区扩容Coordinator发现分区数量变了要给组内成员重新分配。这种重平衡无法避免但它是最温和的一种因为所有消费者都在正常运行。5.2 Rebalance分两步走Revoke与Assign很多文章把Rebalance描述得很神秘其实核心就是两步先收回所有分区再重新分配。具体流程是Coordinator检测到组成员变化通知所有消费者进入“准备重平衡”状态首先执行onPartitionsRevoked回调把当前正在处理的分区交出来然后所有消费者重新发送JoinGroup请求。此时Coordinator会从所有成员中选一个作为Consumer Group LeaderLeader负责用分区分配策略比如RangeAssignor、RoundRobinAssignor、StickyAssignor计算分配方案再通过SyncGroup请求把方案广播给所有成员。这里有个需要记住的点在一个消费者组里第一个加入的Consumer实例会被选为Leader但这个Leader只是个“代表”选举和分配策略都发生在客户端。Coordinator不负责具体分配方案只负责仲裁。5.3 重平衡为什么被称为“stop-the-world”机制Kafka早期版本的Rebalance是全组式的只要组内任何一个成员变化整个组的所有成员会同时停止消费一起参与重平衡然后各自接管新的分区。这在分区数很多、组很大的时候会给业务带来明显抖动。我自己踩过一个印象很深的坑一台消费者机器因为GC停顿导致心跳超时触发RebalanceRebalance期间所有消费者停止消费积压增多积压增多导致其他消费者处理更吃力结果又触发第二轮RebalanceBucket循环整个组在半小时内Rebalance了十几次线上告警响成一片。当时用的还是旧版分配策略全组STW根本躲不开直接舆论。所以在消费者很多比如几十上百个实例的大组里Rebalance风暴是真实存在的灾难。好在新版本里出现了两个缓解机制静态成员Static Membership和增量协作式重平衡Cooperative Rebalance。5.4 减少重平衡的实操手段先说静态成员。只要在Consumer配置里指定了group.instance.id就表示它是“静态成员”。Kafka会通过session.timeout.ms的机制把消费者的离线保护时长拉长很多默认情况下静态成员的离开不会立即触发Rebalance而是等待session.timeout.ms到期让进程重启变得“无感”。这个思路很像把插座式连接改成网线直插——你换设备时网络不能断。然后是增量式重平衡。使用CooperativeStickyAssignor分配策略时Rebalance不再一次性把所有分区全部收回而是只针对有变化的分区做增量调整。没有变化的分区继续由原消费者持有消费不停。这样就实现了“局部重平衡”大幅度减少抖动。我给出的实操清单是心跳、会话与拉取超时参数要一起调session.timeout.ms设成10s以上heartbeat.interval.ms设为session的三分之一比如3smax.poll.interval.ms要根据实际业务耗时设置保守点给到5-10分钟。不要在一个消费者组里塞过多的消费者当消费者数量大于分区数时多余消费者完全空闲纯属浪费资源还增加Rebalance参与方。消费者端逻辑的GC问题一定要监控Full GC导致的心跳超时是Rebalance风暴的头号元凶。用cooperative-sticky分配策略替代默认的range策略特别是组内成员数量大于20时效果立竿见影。生产环境强烈建议开启Static Membership配group.instance.id每次发布代码滚动重启时Rebalance次数肉眼可见地减少。6. 把五个问题串起来一份可直接照抄的生产参数配置与排查路径聊到这儿五大问题都摊开了。但光懂原理不够我最后把整个链路的参数拉一张总表再给你一套排查路径这样你到了现场不至于抓瞎。6.1 关键参数总览不要只记住单个参数要记住组合关系目标需要组合的参数我的配置消息不丢失acksall min.insync.replicas2 unclean.leader.election.enablefalse enable.auto.commitfalse缺一个都不完整生产者可靠retries3 delivery.timeout.ms120000 enable.idempotencetrue重试和时间窗口要匹配有序传输max.in.flight.requests.per.connection1未开幂等时或开启幂等幂等开启后可放宽到5减少重复enable.idempotencetrueBroker层面 消费端唯一键去重/状态机校验消费端幂等是最后防线降低Rebalance风暴session.timeout.ms10000 max.poll.interval.ms300000 group.instance.idxxx partition.assignment.strategycooperative-sticky这几个参数缺一不可注意max.poll.interval.ms和max.poll.records这两个参数兄弟配合使用如果你处理一条消息平均要5秒单次pull又拉回500条那处理完一轮就需要2500秒早就超过5分钟上限必然触发Rebalance。正确做法是调低max.poll.records比如每次只拉100条让处理一轮的耗时压在max.poll.interval.ms之内。6.2 遇到问题后按链路排查如果线上出现消息丢失别急着改参数按这个顺序来查查Producer日志有没有NotEnoughReplicasException、超时重试、序列化异常生产者的错误日志是一切线索的起点。查Broker日志有没有LeaderAndIsr变化记录有没有unclean leader election日志查Consumer日志有没有提交Offset的时间点有没有OffsetOutOfRangeException查监控ISR有没有收缩过UnderReplicatedPartitions指标有没有报警这些监控项要在问题出现前就配好。如果是顺序乱了先确认是不是“同一个Key的消息进了不同分区”可以用一条命令查消息的partition归属kafkacat或kafka-console-consumer加--property print.partitiontrue再看消费者是不是多线程乱掉了。如果是重复消费先看是不是Rebalance导致的重复重点是消费者有没有在处理完前提交Offset再看业务侧有没有幂等兜底。所有机制都挡不住第一类重复只要Consumers数和Partitions数没对上Rebalance一发生重复就是必然的。6.3 我的一点私人体会这么多年用下来我最大的体会是Kafka的参数不是越“保险”越好而是要和你的业务代码咬合。很多人在生产环境把可靠性相关的参数全部拉满结果吞吐掉得厉害又改回默认值改回来之后丢消息了又怪Kafka。其实问题不在参数本身在于你根本没有梳理过自己消息链路里的真实风险点。比如订单类业务丢一条消息可能就是资损但日志分析类业务丢个千分之一真无关紧要。与其追求“绝对可靠”不如把关键链路的可靠性做扎实把非关键链路的吞吐拉满。先定业务目标再动参数不要上来就复制一篇“最佳实践”就完事——那样你连自己到底保护了什么都不知道。排查问题的时候我习惯把所有关键配置打印在服务启动日志里线上出事故时能一眼看到当前跑的是哪套参数组合这一点看起来不起眼但关键时刻能帮你省下半小时。