
但凡是用过 Kafka 消费端的人大概率都经历过这种场景消费者进程明明没挂日志里却突然出现一堆 “rebalance” 相关告警紧接着消费任务停顿几十秒恢复之后消息又重复处理了好几轮。我在生产环境里第一次踩到 rebalance 的坑就是因为发布新版本时滚动重启消费实例结果整个消费者组像是被按住了暂停键下单接口被重复触发了十几次排查到半夜结论就是一句话rebalance 风暴。那之后我花了不少时间把 Kafka 的 rebalance 机制从头到尾研究了一遍。这个机制可以说决定了消费端的高可用和数据一致性但它藏得很深平时跑得稳的时候你根本感觉不到它的存在一旦出了问题它就是整个链路里最让人头疼的那一环。这篇文章我想用实战视角把 rebalance 机制讲透它为什么存在、触发后底层到底做了什么、哪些参数能救你命、哪些错误日志意味着什么以及我踩过的一些坑。1. Rebalance 到底在解决什么问题1.1 消费者组负载均衡的最小单元你要理解 rebalance先得理解消费者组。Kafka 的消费端不是单打独斗的同一个group.id下的所有消费者实例会组成一个逻辑上的“消费组”。一个主题的分区会被这个组里的成员瓜分每个分区在同一个时刻只能被组内的一个消费者实例消费。打个比方你开了一家奶茶店有 8 个出餐窗口分区店里有 3 个员工消费者实例。每个窗口同一时刻只能有一个人负责接单这 3 个人需要把 8 个窗口划分清楚你管哪几个、我管哪几个。这种“重新划分窗口归属”的动作落到 Kafka 里就叫 rebalance。它解决的核心问题有三个故障转移某个消费者实例宕机了它手里的分区不能留在那里没人管需要分给其他存活成员。弹性扩缩容消费能力不够了你多加几个消费者实例分区需要重新分配反过来缩容了也一样。分区数量变化主题从 8 个分区扩容到 16 个分区原本的分配关系也必须重新算一遍。没有 rebalance消费端就无法动态响应这些变化。但你也会发现rebalance 本质上是一次“全体成员重新分配”的动作执行过程中必然会造成消费短暂中断这也就是后面所有麻烦的根源。1.2 什么情况会触发 rebalance我梳理了一下实际生产中常见的触发场景基本可以分成四类。第一类是成员变化。组内有消费者实例加入或者退出。退出又分主动和被动主动是进程调用close()或者在优雅关闭时发起了 LeaveGroup 请求被动是消费者心跳超时被协调器判定为“死亡”然后被强制踢出组。第二类是订阅变化。消费者订阅的主题集合发生了变化。比如原来订阅topic-a后来又加订阅了topic-b或者订阅用了正则表达式原先匹配 2 个主题后来新增了一个匹配主题。这都会触发整个组重新评估订阅关系。第三类是分区变化。主题的分区数增加比如运维对某个大流量主题做了扩容把分区从 12 个扩到 24 个。这种情况下消费者组成员不变但分区归属必须重算。第四类是分配策略不一致。组内成员配置的partition.assignment.strategy不一样协调器无法兼容也会反复触发 rebalance。这个问题比较隐蔽多见于组内混合部署了不同版本的客户端或者有人手滑改了配置没同步。很多新手容易有一个误解以为只有“消费端”发生变化才会触发 rebalance。实际上主题分区扩容这种“生产端数据源侧”的变化同样会触发而且往往发生在流量高峰影响面更大。1.3 Rebalance 的“停市”特性rebalance 期间整个消费组会处于一个不稳定状态分区和消费者的归属关系是不确定的。在传统的 eager rebalance 协议下甚至会出现“先把所有分区从所有消费者手里收回再重新分配”的过程这就等于所有消费者在这一瞬间集体停摆。这就像奶茶店的 3 个员工正在各忙各的窗口突然店长喊了一句“所有人放下手里的单子重新分一下窗口”。那么一瞬间所有窗口都没人接单。如果这个“重新分窗口”的过程反复发生店里就是一团乱麻。这也是很多人对 rebalance 又爱又恨的原因。爱是因为没有它就无法动态容错恨是因为它在某些场景下会把消费集群搞成“抖动状态”消息延迟、重复消费全来了。2. 核心机制拆解从 JoinGroup 到 SyncGroup2.1 协调者组的管理中枢rebalance 不是客户端之间自己商量出来的而是由一个“裁判”来安排的这个裁判就是 GroupCoordinator也就是组协调者。每个消费者组都会对应一个协调者它并不是独立的进程而是某个 Kafka broker 上的一个逻辑角色。具体是哪个 broker 来当协调者是根据消费者组的group.id算出哈希值再对__consumer_offsets的分区数量取模找到对应的 partition然后由这个 partition 的 leader 副本所在的 broker 来担任协调者。协调者负责管理组的元数据组成员列表、每个成员的订阅信息、分区分配的方案以及各个成员上报的消费位移。这里有个容易踩的坑如果你的__consumer_offsets这个内部主题的副本因子配得比较低有些测试环境甚至只有 1那协调者所在的 broker 一挂整个消费者组都要经历一轮重新发现协调者的过程期间消费自然也会卡顿。2.2 完整的 rebalance 流程一次 rebalance 大致有四个阶段。第一阶段查找协调者。消费者通过FindCoordinator请求找到当前组对应的协调者 broker。这个阶段一般在消费者启动或者协调者发生变化时触发。第二阶段加入组也就是 JoinGroup。消费者向协调者发送 JoinGroup 请求带上自己的成员 ID、订阅信息以及支持的分配策略列表。协调者会从当前存活的成员中选出一个作为 leader。这个 leader 并不是一个特殊角色它只是一个“代表”负责帮所有人计算分区分配方案。选 leader 的方式很简单谁先来或者谁的成员 ID 排在最前不同版本细节略有差异但都是协调者内部的固定逻辑。第三阶段同步组状态也就是 SyncGroup。leader 消费者根据选定的分配策略算出每个成员应该拿到哪些分区然后把结果通过 SyncGroup 请求发给协调者。协调者再把最终分配方案分发给组内每一个成员。第四阶段稳定状态。所有消费者拿到属于自己的分区分配结果后rebalance 结束进入正常的心跳维持和消息拉取阶段。这里有一个非常关键的细节在第二阶段协调者有一个“等待时间”它会等一段时间来收集尽可能多的成员加入请求。如果某些成员迟迟不申请加入协调者也不会无限等下去超时后就直接用当前已有的成员先进行分配。这一等往往就是几秒甚至更久于是你会看到消费组整体“停顿”几秒的现象。2.3 组状态机与分代机制协调者内部管理每个组的生命周期是在一组状态之间来回切换的大致包括Empty组内没有成员或者是组刚建出来。PreparingRebalance已经发现有成员加入、离开或者订阅变化准备开始新一轮 rebalance。CompletingRebalance正在等 leader 计算分配并等待各位成员回复 SyncGroup。Stable分配完成组成员处于稳定消费中。Dead组被彻底移除。理解这个状态机的价值在于排障。比如你用工具查看组的当前状态发现它一直停在PreparingRebalance或者CompletingRebalance那基本上可以判断有成员反复地“加入又退出”形成了 rebalance 抖动。分代机制也很重要。每开始一次 rebalance代的编号就会加一。你在代码里调用commitSync()的时候请求里会带上当前的 generation。如果提交位移的时刻已经开启了新一轮 rebalancegeneration 对不上协调者就会拒绝这次提交客户端会抛出CommitFailedException。这是很多“重复消费”问题的直接原因你以为消息已经处理完了、位移也提交了实际上因为 rebalance 导致提交失败重启后又要从头消费。2.4 Eager 与 Cooperative 的区别如果按协议版本分rebalance 有两种大风格理解它们的区别对排障特别重要。Eager rebalance 是旧版协议的做法核心特征是“全量重分配”。每次进入 rebalance所有消费者都必须先把当前手里持有的分区全部释放掉然后重新加入组等待分配。过程简单粗暴但有个致命弱点瞬间全组停摆并且停摆期间所有分区都没人消费。Cooperative rebalance 是 KIP-429 引入的增量式变更。它的核心思想是“能不动就不动”只有受变化影响的分区才被重新分配保持稳定的分区继续由原来的消费者持有。这样可以尽量减少 rebalance 对消费吞吐的影响。不过兼容性上要留意同一消费者组内所有成员必须使用相同的协议版本。如果新老客户端混用很容易出现协议协商失败的问题。我建议在升级客户端时统一所有实例的版本尽量别搞“先升一半观察一下”这种策略。3. 分区分配策略不是玄学是算法3.1 四种常用的分配器rebalance 过程中分区具体怎么分取决于每个消费者客户端配置的partition.assignment.strategy。Kafka 自带了几种分配器我强烈建议你搞清楚它们之间的差别因为选错了是真的会“吃不均”。RangeAssignor也叫范围分配器。它会把每个主题的分区按照消费者的字典序排序然后一段一段地分配给各个消费者。比如一个主题有 6 个分区组里有 2 个消费者分区 0-2 给第一个消费者3-5 给第二个消费者。它的缺点是如果组内消费者订阅了多个主题而且各主题分区数不能均匀整除消费者数靠前的消费者就会多分到分区很容易出现“旱的旱死涝的涝死”。RoundRobinAssignor也叫轮询分配器。它把所有订阅主题的分区全部列出来然后按消费者轮流分配。这样各个消费者持有的分区数量会比较均衡但连续性和局部性较差。StickyAssignor粘性分配器。它在保证均衡的前提下尽量保留上一次分配的结果。这能减少 rebalance 之后分区在不同消费者之间“搬家”的数量也就减少了重复消费和状态重建的成本。CooperativeStickyAssignor协作粘性分配器。它是和 Cooperative rebalance 协议配套的分配器既保证均匀又尽量保留已有分配同时支持增量式变更。目前新版本里比较推荐这个。我用一个简单场景对比一下假设一个组里有消费者 C0 和 C1订阅了主题 T04 个分区和主题 T12 个分区。Range 分配器很容易让某个消费者分到 4 个分区另一个只分到 2 个RoundRobin 则可能让两边各拿 3 个。差别就是这么大。3.2 多策略配置的坑有些团队会在配置里写多个策略比如partition.assignment.strategyorg.apache.kafka.clients.consumer.RangeAssignor,org.apache.kafka.clients.consumer.CooperativeStickyAssignor。这个配置的意思是消费者支持这两种策略由协调者从组内成员的策略列表中进行交集选择。听起来灵活但实际很容易埋雷。如果组内一部分成员只支持 A 策略另一部分只支持 B 策略交集为空协调者就无法统一分配方案rebalance 会反复失败。特别是刚升级客户端协议版本时最容易出这种问题。我的建议是组内所有成员用同一种分配器并且显式配置不要隐式依赖默认值。生产环境我现在基本固定用CooperativeStickyAssignor除非对端还有一些老版本客户端混跑。4. 实操配置与代码防线4.1 三个超时参数组成了消费端的生命线第一个是session.timeout.ms。它用来检测消费者进程是否“存活”。如果协调者在这么长的时间内没有收到消费者的心跳就认为该消费者已经挂了会把它踢出组并触发 rebalance。值太小网络抖动稍微大一点就会被误杀值太大消费者真挂了也要等很久才触发转移故障恢复时间变长。第二个是heartbeat.interval.ms。它决定消费者发送心跳的频率。心跳的职责是告诉协调者“我还活着”。一般建议设置为session.timeout.ms的三分之一左右。比如session.timeout.ms设 30 秒心跳间隔就设 10 秒。这样即使一个心跳包偶尔丢了协调者也能在 session 超时之前收到下一次心跳。第三个是max.poll.interval.ms。它决定两次poll()调用之间的最大间隔。如果消费者拉取消息后处理时间太长超过这个间隔没有发起下一次 poll协调者会认为这个消费者“卡死了”强制把它移除并触发 rebalance。这个参数很多人容易忽略但它恰恰是 rebalance 风暴最常见的原因。下面是一个我常用的消费端基础配置示例Properties props new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, broker-1:9092,broker-2:9092); props.put(ConsumerConfig.GROUP_ID_CONFIG, order-group); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, org.apache.kafka.common.serialization.StringDeserializer); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, org.apache.kafka.common.serialization.StringDeserializer); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 200); props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 300000); props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 30000); props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, 10000); props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG, org.apache.kafka.clients.consumer.CooperativeStickyAssignor);4.2 消费主循环如何优雅地处理消息和位移我见过很多 rebalance 问题本质上不是机制的问题而是消费主循环写得不对。最常见的错误是在while循环里拉取消息后直接同步处理业务逻辑处理完再 commit。如果一个批次的记录处理时间超过max.poll.interval.ms就触发了 rebalancerebalance 又导致 CommitFailedException兜底逻辑又去重试结果越跑越乱。要做到稳定我的几个经验是使用enable.auto.commitfalse自己控制提交时机避免自动提交与 rebalance 重合造成不可控的重复消费。把消息处理做成异步批处理或者把处理的活儿丢给独立线程池poll 线程只负责拉取和调度保证 poll 能按时被调用。调整max.poll.records让它和你单批处理耗时匹配。同样的处理逻辑一次拉 500 条处理不过来就降到 100 条、50 条稳定压倒一切。在ConsumerRebalanceListener里做两件事分区被收回时同步提交一次位移并清理本地状态分区被分配时再重建状态比如重新定位本地处理游标。下面是一个经典的监听器写法consumer.subscribe(Collections.singletonList(order-topic), new ConsumerRebalanceListener() { Override public void onPartitionsRevoked(CollectionTopicPartition partitions) { // 在分区被回收前把已经处理完的位移提交掉尽量避免重复消费 consumer.commitSync(); // 同时清掉本地缓存、聚合值等状态 localCache.clear(); } Override public void onPartitionsAssigned(CollectionTopicPartition partitions) { // 重新初始化本地状态 for (TopicPartition partition : partitions) { Long offset localOffsetStore.get(partition); if (offset ! null) { consumer.seek(partition, offset); } } } });很多人只重写了onPartitionsAssigned忘了onPartitionsRevoked。这是非常危险的。因为 rebalance 发生时你还没处理完一批消息旧分区被分配给别人之后别人会从已提交的位移开始消费。如果你没有在 revoked 阶段提交位移可能还停在几批之前对端一接手就是大量重复消息。4.3 静态成员减少不必要的 rebalanceKafka 从 2.3 版本开始支持静态成员机制。简单说给消费者实例配置一个固定的group.instance.id这样每次重启时协调者能认出“还是原来那个成员”而不会因为进程重启就触发一次完整的 rebalance。这在滚动发布场景里特别有用。没有静态成员时发布一次就触发一次 rebalance分区在成员之间来回换主人带来大量无效的重复消费。有了静态成员只要在 session 超时时间内重启完成协调者会保留这个成员的身份分区归属几乎不动。当然静态成员也有它的代价如果一个实例直接宕机且短时间内起不来那么到 session 超时之前它的分区一直是无人消费的状态因为协调者还在等它“回归”。所以这个功能适合配合优雅停机使用不适合暴力 kill。5. 常见问题与排查实录5.1 现象一日志里反复出现 CommitFailedException这类错误的典型日志是Commit cannot be completed since the group has already rebalanced and assigned the partitions to another member.意思很明确你在提交位移时组已经重新平衡了分区也不归你管了所以提交失败。如果你没捕获这个异常默认行为是继续消费已经拉取到的记录但位移不会更新。重启之后这些消息会被再次消费。我的排查思路一般是三步第一步先看是不是max.poll.interval.ms太小拖累了批处理速度。第二步看业务代码里有没有在 poll 循环内做耗时操作比如调远程接口、写日志、批量写数据库。第三步看提交方式。如果是同步提交确认提交前有没有在 revoked 阶段处理过。如果是异步提交要处理好回调里重试逻辑别引入乱序提交。5.2 现象二rebalance 风暴rebalance 风暴是影响最严重的一种情况。表现是消费者实例反复地退出组又立刻重新加入组状态一直切不到 Stable整个消费组都在空转消息延迟飙升甚至有消息持续无人消费。我遇到过一个典型的 case某台消费实例部署的机器 GC 停顿比较长导致心跳发送延迟被协调者判定为 session 超时踢出了组。踢出后触发 rebalanceGC 又恢复过来这个实例马上重新入组再次触发 rebalance。就这样来来回回地抖一次发布引发了上百次组状态变更。要规避这种问题除了调大session.timeout.ms和适度增加heartbeat.interval.ms还要关注消费机器的 GC 停顿、网络抖动这些底层稳定性问题。参数能缓解症状但根本还是要让消费实例本身足够健壮。5.3 现象三消费者数量大于分区数量有时候你以为多加实例就能提升消费速度但如果分区数是固定的 8你却起了 10 个消费者实例那么最终只会由 8 个实例分到分区多出来的 2 个实例会一直空转。它们加入组时会触发 rebalance但由于没有多余的分区可分分配完成后依然没有活干。这不是 bug而是机制本身的特性。但有个坏处多出来的实例会时不时导致组内成员变化增加不必要的 rebalance。所以起消费者实例之前想清楚分区数到底是多少一般建议消费者实例数和分区数保持一致即可不需要盲目地加。5.4 排障命令和监控指标遇到 rebalance 问题通常我会先用命令行看消费者组的当前状态。以 3 节点集群为例bin/kafka-consumer-groups.sh --bootstrap-server kafka1:9092,kafka2:9092,kafka3:9092 \ --describe --group order-group输出里会显示当前消费者组的状态是Stable还是PreparingRebalance。如果是Stable再进一步看每个分区对应的CURRENT-OFFSET和LOG-END-OFFSET算一下 lag 是否正常。如果状态始终不稳定基本就能锁定 rebalance 在持续发生。监控指标方面下面几个值得重点关注kafka.consumer:typeconsumer-coordinator-metrics下的commit-latency-avg、heartbeat-rate、heartbeat-response-time-max。kafka.consumer:typeconsumer-metrics下的records-lag-max用来观察消费延迟是否在扩大。处于PreparingRebalance状态的时长和次数需要靠 broker 侧的指标和日志综合判断比如 broker 日志里频繁出现的组状态转换记录。结合这些指标再回头调参才不会只是“拍脑袋试”。我自己就经历过一次因为只调max.poll.records却忽略了max.poll.interval.ms的无效操作折腾了大半天才意识到两个参数是连锁影响的。6. 我的一些实操体会如果你现在正被 rebalance 问题困扰我建议你先做一件事把enable.auto.commit改成false把max.poll.interval.ms放宽到至少比你处理一批消息所需最长时间再大 50%然后观察组状态能不能稳定下来。很多时候这一步就能解决七八成的问题。另一个建议是在测试环境里故意制造一次消费者宕机和重启观察消费者组的 rebalance 耗时、重复消费量。你会发现哪些代码路径在 rebalance 期间是不安全的比如本地缓存没清理、位移没提交、数据库幂等没做好。这些内容在正常跑的时候完全看不出来只有在 rebalance 的那一刻才暴露。Kafka 的 rebalance 机制本身并不花哨但它是消费端一致性的一道闸门。理解它的触发条件、流程和参数之间的关系比背一百个调优参数都管用。希望这篇内容能帮你少踩几个我踩过的坑。