ARTICLE DETAIL

资讯详情

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

Kafka消息堆积百万条?一次rebalance风暴引发的消费组故障复盘

Kafka消息堆积百万条?一次rebalance风暴引发的消费组故障复盘 凌晨两点报警群里蹦出一条消息Kafka 消费组 lag 突破了 100 万条。说实话第一眼看到这个数字我第一反应是监控脚本是不是写错了位数。确认了三遍之后数字是真的而且还在涨一分钟涨一万多。那时候脑中闪过的第一个念头是完了又得加消费者实例了。但后来真正解决问题的时候别说加机器连并发都没调只是把消费者配置里一个默认参数改了。参数生效之后堆积量以肉眼可见的速度往下掉半小时清掉了接近 100 万条。这篇文章就把那次从错判到救火的完整过程复盘一遍包括排查命令、触发机制的原理、参数调整的具体操作以及防止复发的消费侧改造。内容面向正在做实时数据链路、被 Kafka 消息堆积困扰的开发和运维同学。就算你用的是 Java 原生客户端而不是 Spring Kafka这篇的思路和参数同样适用。1. 先别急着清堆积从 lag 分布里读出真实线索1.1 查消费组状态先回答两个问题遇到消息堆积第一个动作永远是查消费组状态而不是拍脑袋调参数。我最常用的命令是kafka-consumer-groups.shkafka-consumer-groups.sh \ --bootstrap-server kafka01:9092 \ --group order-consumer \ --describe跑出来的结果大概是这个样子GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST order-consumer order-topic 0 24500000 25400000 900000 consumer-1-xxx /10.0.0.12 order-consumer order-topic 1 18000000 18005000 5000 consumer-2-xxx /10.0.0.13 order-consumer order-topic 2 15200000 15200000 0 consumer-3-xxx /10.0.0.14这张表的信息量比我一开始想的要大得多。CURRENT-OFFSET 表示消费组在这个分区上已经提交的位移LOG-END-OFFSET 表示 broker 端这个分区的最新位移两者相减就是 LAG。很多人只看 LAG 大小但忽略了另一个关键维度看 LAG 是集中在个别分区还是均匀分布在所有分区。差别非常大前者指向单个消费者实例后者通常说明整个消费组的整体吞吐不足。拿到输出后我还会过一遍 CONSUMER-ID 和 HOST 列。如果这一列里出现多个实例且分布均匀说明消费者数量是够的如果某个分区的 CONSUMER-ID 为空说明这个分区正处于无人消费的状态往往就是 rebalance 正在发生或者消费者已经挂掉。1.2 连续跑两次 describe判断消费是慢还是卡死只查一次还不够。我当时的做法是隔 30 秒再跑一次同样的命令对比两次的 CURRENT-OFFSET。如果 CURRENT-OFFSET 在递增说明消费线程还在跑只是速度跟不上生产速度这种场景的思路应该是扩容或者优化单条消息的处理逻辑但如果 CURRENT-OFFSET 完全不动问题就没这么简单了消费者大概率不是慢而是根本没在消费。那次事故里我把两次命令的输出放在一起对比0 号分区的 CURRENT-OFFSET 45 分钟都没动过而 1 号、2 号分区的位移每隔几秒都在往前走。再结合堆积量 90 万条集中在 0 号分区几乎可以锁定嫌疑目标负责 0 号分区的那个消费者实例出了状况。到这里为止一切都还很常规我甚至已经准备联系 DBA 扩容消费组了。但接下来翻消费者日志时看到的东西让我把话咽了回去。2. 堆积元凶不是消费速度而是 rebalance 风暴2.1 消费者被踢出的机制心跳与 poll 是两个完全不同的维度Kafka 判断一个消费者是否存活实际依赖两个彼此独立的维度心跳线程和 poll 调用。每个消费者内部都有一个独立的心跳线程定期向协调器发送心跳包证明我还活着。同时业务线程会反复调用poll()这个操作既是拉取消息也是在告诉 Kafka我还在正常处理数据。这里有两个非常容易混淆的参数session.timeout.ms和max.poll.interval.ms。前者管心跳默认 10000 毫秒也就是说消费者如果在 10 秒内没有心跳协调器就会判定它死亡把它的分区分给别人。后者管 poll 的频率默认 300000 毫秒也就是 5 分钟如果消费者的两次 poll 调用之间超过 5 分钟即使心跳一直正常协调器也会认定它是一个活着的僵尸——明明进程还在但已经不干活了于是强制把它从消费组里移除触发 rebalance。用生活化的话说session.timeout 负责判断人有没有断气max.poll.interval 负责判断人是不是在摸鱼。之前我们总盯着 session.timeout 看但真正容易踩雷的其实是 max.poll.interval.ms因为它的默认值只有 5 分钟而一个下游调用超时随便就能超过 5 分钟。2.2 Rebalance 风暴的滚雪球效应那次事故的现场我当时用kafka-consumer-groups.sh查完 lag 之后去消费者应用日志里翻发现满屏都是类似的日志[WARN] org.apache.kafka.clients.consumer.internals.ConsumerCoordinator: The consumer has timed out, meaning that it has not polled for 300000ms.再往下翻还能看到这个消费者反复加入消费组的记录。这一切串起来真相就浮出水面了负责 0 号分区的消费者在处理某条业务消息时需要调用一个下游服务的接口结果这个接口不争气一卡就是 12 分钟。处理线程被下游接口阻塞自然没机会回到 poll 循环里继续调用 poll。5 分钟的max.poll.interval.ms一到协调器立刻把它踢出消费组触发 rebalance。Rebalance 本身是一个全局停止消费的过程。在再平衡期间消费组里的所有消费者都会暂停消费等待分区重新分配。好巧不巧其他消费者拿到 0 号分区后同样会碰到那批需要调用下游接口的积压消息于是也卡住也被踢出再度触发 rebalance。整个消费组陷入了踢出 - 再平衡 - 接手 - 再踢出的恶性循环。我后来计算过这 45 分钟里消费组真正在消费数据的时间可能不到 5 分钟其余时间全耗在 rebalance 上。堆积量就是这么滚起来的消费速率不是变慢了而是直接变成了零生产端每秒钟还在往 partition 0 里写几百条数据lag 自然像雪球一样越滚越大。搞清楚这一点之后我彻底放弃了加机器的念头——加再多的消费者实例只要处理逻辑里的阻塞区不解决新实例加入之后照样会在 5 分钟后被踢出。3. 救火操作一个参数让堆积半小时清空3.1 参数调整方案max.poll.interval.ms 与 max.poll.records 的组合既然瓶颈是下游接口阻塞导致 poll 超时被踢那最直接的办法就是把消费者的容忍时间拉长。Java 原生客户端里关键参数是max.poll.interval.ms我从默认值 300000 毫秒直接调到了 1800000 毫秒也就是 30 分钟。同时把max.poll.records从默认的 500 调到了 200避免单轮 poll 拿太多消息导致总处理时长过长。原生客户端这样改Properties props new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafka01:9092,kafka02:9092); props.put(ConsumerConfig.GROUP_ID_CONFIG, order-consumer); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); // 核心把两次 poll 之间的最长间隔从默认 5 分钟调大到 30 分钟 props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 30 * 60 * 1000); // 辅助降低单轮拉取的记录数控制单轮处理耗时 props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 200); KafkaConsumerString, String consumer new KafkaConsumer(props);如果是 Spring Kafka 项目配置更简单在 application.yml 里加两行即可spring: kafka: consumer: max-poll-records: 200 properties: max.poll.interval.ms: 1800000这里要特别说明一下max.poll.interval.ms调大本质上是给业务处理争取更长的喘息时间让消费者别被轻易踢出消费组而max.poll.records调小则是从源头控制单轮处理量。两者是配合关系不是替代关系。如果只调大前者遇到真正致命的处理异常时消费者照样会卡住只是要等到 30 分钟后才被发现如果只调小后者下游接口一旦阻塞超过 5 分钟照样触发超时踢出。3.2 为什么看似瞬间清空背后发生了什么参数调整完成后重启了消费组。重启后的一两分钟内能看到日志里还有一两行 rebalance 的痕迹这是消费组重新给自己分配分区的正常过程。之后整个消费组就安静下来了没有新的踢出记录没有新的 rebalance。随后我盯了一道监控lag 数据是这样走的时间LAG状态02:101,000,000调整参数并重启消费者02:15750,000消费组稳定无 rebalance 日志02:20400,000消费速率恢复到正常水平02:355,000堆积接近清空02:400完全消化看到 02:35 那个数字时我确实有瞬间清空的错觉。但冷静下来复盘这个瞬间背后其实有一笔很容易算的账消费者的真实处理能力一直是每秒几千条真正的问题只是它每隔 5 分钟就被踢出一次根本没有连续消费的机会。阻塞的下游接口在 20 分钟后自己恢复了poll 超时也不再发生消费线程终于可以一口气把所有积压数据消费完。也就是说参数只是把错误的死亡判定纠正了消费能力从头到尾都没坏过。我也想过要不要把max.poll.interval.ms直接调成远超下游接口超时时间的值这里其实有一个权衡调太大如果消费者真的遇到死循环或者内存溢出协调器需要很久才能把它的分区转移给其他消费者故障恢复时间会变长。我选择 30 分钟是因为当时下游接口的超时上限设在 15 分钟30 分钟留了一倍的余量。如果下次下游超时设计改了这个参数还是要跟着重新评估。4. 根治方案把轮询与业务处理解耦再上监控4.1 单线程 poll 线程池异步处理的正确姿势调参数解决了眼前的堆积但我在这里明确提醒一句max.poll.interval.ms调大只能算治标业务处理里的阻塞调用不拆掉下次换个场景照样出事。最稳的消费模型应该是poll 循环里只做最轻量的解析和投递把真正耗时的下游调用、数据库写入等操作丢给独立的线程池。很多同学会在这里踩一个坑直接在业务线程池里消费ConsumerRecords然后还在主线程里提交位移。KafkaConsumer 本身不是线程安全的如果多线程共享同一个 consumer 实例去 poll 或者 commit会直接抛异常或者出现诡异的位移漂移。比较稳妥的做法是主循环只负责 poll拿到数据后马上把 records 交给有界队列和线程池位移提交单独用一个线程或者使用commitAsync同时保证只有当一批消息处理完成之后才提交该批位移。如果确实需要在多个线程里同时消费数据我的建议是多分区 - 多 consumer 实例而不是让多个线程共享一个 consumer也就是每个线程持有自己的 KafkaConsumer各消费各的分区。这样既绕开了线程安全问题又天然把并发度撑起来了。代价是线程数直接等于分区数而且要把 offset 持久化工作做好。4.2 Spring Kafka 场景下的配置与坑如果你用 Spring Kafka事情看起来简单但坑反而更隐蔽。KafkaListener默认是单消费者单线程模式listener 方法里执行的业务逻辑越长poll 循环被阻塞的时间就越长最终同样会顶到max.poll.interval.ms。很多人以为只要把 listener 的 concurrency 调大就能解决结果发现调大之后依然跟着报警原因是每个并发消费者线程都有各自的 poll 超时。在 Spring Kafka 里除了设置max.poll.interval.ms之外还要注意 AckMode。如果用的是默认的BATCH模式那么 listener 处理完一批消息之后返回时才算 ack处理过程中一旦抛异常这批消息可能会被重新拉取形成死循环。对于追求吞吐的实时链路我通常建议用MANUAL_IMMEDIATE在处理逻辑中明确调用acknowledgment.acknowledge()这样位移可以做得更细失败重试时也不会把整批消息再拉一遍。还有一个很多人不知道的细节KafkaListener如果处理消息的时间超过了max.poll.interval.msSpring Kafka 会打印一条包含Consumer ... has not polled for ...的 WARN然后消费者会被强制停止容器发起重新定位。这个日志特征非常容易被误判成网络问题所以一看到它第一反应应该是我的 listener 是不是卡太久了。4.3 lag、rebalance 与已提交位移的三层监控治本的另一半是监控。我现在的经验是光盯着 lag 一个指标远远不够。Lag 变大时必须能立刻区分两种情况消费组在正常消费但速度跟不上还是消费组已经被踢得只剩骨架。正常情况下比较两次提交位移如果 CURRENT-OFFSET 在稳步前进说明消费还在跑只是慢如果 CURRENT-OFFSET 纹丝不动几乎可以断定消费者被卡住或者已经因为 poll 超时被移出。所以我在监控体系里做了三层指标第一层是 lag 总量和 lag 增长速度阈值一般定在该消费组正常 lag 基线的 5 倍左右防止小波动误报第二层是消费组成员数量变化和 rebalance 次数通过 JMX 的消费者协调器指标可以拿到一旦在短时间内出现多次 rebalance基本就是风暴苗头第三层是已提交位移的增长速率这个指标被很多人忽略了但恰恰是区分慢和死的最有力证据。三层指标配合起来报警触发之后能直接给出初步结论比如lag 高但位移在增长建议扩容或者lag 高位移不动且 rebalance 次数激增建议查 poll 超时和业务阻塞。这套逻辑确定之后再遇到堆积就不需要我凌晨爬起来对着 describe 命令发呆纠结了。5. 复盘时我总结的三条血泪法则5.1 三条判断法则那次事故之后我把排查 Kafka 消息堆积的流程压缩成了三条判断法则分享出来供你参考。第一条堆积大不等于消费慢。看到 lag 高先看消费组里消费者数量是否稳定rebalance 是否频繁。如果消费者数量在 10 分钟内反复变化那问题大概率出在稳定性上而不是吞吐量上。第二条max.poll.interval.ms的默认值是为正常的轻量业务逻辑设计的任何单条消息处理可能超过 5 分钟的场景都要提前把这个参数纳入评估范围。不要等到报警了才想起有这个东西。第三条判断消费进程是死是活别只看进程状态和心跳要看 CURRENT-OFFSET 是不是在动。这是最快、也最不容易被表象骗到的判断方式。5.2 为什么调参而不是加机器最后说说为什么当时我选择调参数而不是直接加机器。从容器里加两个 pod 看起来更快但仔细想想如果根因是max.poll.interval.ms导致的 rebalance 风暴新加入的消费者会被协调器重新分配分区。它一接手就开始消费积压消息然后同样碰到那个需要调用下游接口的阻塞逻辑一样会超过 5 分钟的 poll 超时阈值。最终结果是新消费者加入后忙活几分钟然后同样被踢出消费组进一步动荡堆积只会更严重。调参看似简单其实是在恢复消费组最稀缺的资源——稳定性。一个稳定的消费组即使消费速度不是特别快也能持续地把 lag 降下来一个反复 rebalance 的消费组哪怕每个消费者单机吞吐都不错整体能用吞吐也趋近于零。这个道理放到其他分布式系统里也一样成员频繁变动的集群往往比稳定的慢集群更容易出灾难性问题。那次之后我把max.poll.interval.ms纳入了一轮所有消费组盘的改造清单里凡是业务逻辑中存在外部调用或者重计算的地方都单独评估过这个参数。只要 poll 循环和业务逻辑足够分离这个参数未来大概率只会在配置里躺平但每一次堆积报警背后它都可能是第一个需要被想起的名字。
返回列表