跳至主要內容

十万个why:消费者线程明明还在正常打日志,为什么 Kafka 还是判定它“死亡”并触发了 Rebalance?

程序员小富大约 5 分钟

大家好,我是小富~

《十万个why》系列持续更新中,这个系列都是我面试过程中接触过的问题,也是我之前工作踩过的坑。

线上发版或者业务高峰期,有时会遇到一个奇怪的现象:

消费端的服务进程没挂,JVM 内存正常,日志也在正常输出,显示业务代码还在继续处理数据。但 Kafka 服务端却判定该消费者已失效,强行将其踢出消费组,并触发了 Rebalance。

这种在进程和日志均正常的情况下依然被判定为失效的现象,在技术上通常被称为假死

为什么有心跳也会被判定死亡?

Kafka 极早期的版本,消费端是单线程的。

这个单线程要同时负责两件事:拉取数据并执行写库、打日志等业务逻辑,还要给服务端发心跳包,报告自己还活着。

但这个设计在性能上存在缺陷:如果业务逻辑卡住了,比如查数据库超时,或者 JVM 发生 Full GC,主线程被卡住,就无法按时发心跳。服务端等不到心跳,就会误认为客户端挂了,立刻触发 Rebalance。

后来 Kafka 把拉数据和发心跳彻底拆成了两个线程。

双线程架构

重构之后,启动一个 Kafka 消费者,底层实际上是两个线程在协作。

心跳线程只负责在后台按照固定频率给 Broker 发心跳包。只要心跳不断,服务端就认为你的进程还活着。它判断存活的阈值是 session.timeout.ms,默认是 45 秒。

另一个业务主线程,负责在循环里调用 poll() 拉取数据并执行业务逻辑。它判断是否卡死的阈值是 max.poll.interval.ms,默认是 5 分钟。

这里有个最容易发生冲突的地方。

假设你的消费者一次性拉了 500 条消息,因为下游写库慢,主业务线程处理这 500 条消息一共花了 6 分钟。

在这 6 分钟里,主线程正在处理数据并输出日志。由于这 500 条数据还没处理完,它是没办法回到循环去调用下一次 poll() 的。

可后台的心跳线程依然在按时给服务端发心跳。但在服务端看来,距离这台机器上一次调用 poll() 已经过去 6 分钟了,超过了规定的 5 分钟死线。

服务端就会判定,这个消费者虽然还有心跳,但它已经失去了消费能力(占着分区却不拉取新数据,导致数据积压)。为了不影响整个消费组的吞吐,必须判定它为假死,强行将其踢出消费组,这就引发了 Rebalance。

怎么避免假死?

知道是由于业务处理太慢,导致 poll() 间隔超时,其实解决思路非常清晰。

1. 减小拉取批次

客户端配置里,限制每次 poll() 拉取的最大消息数:

max.poll.records: 50

把大批次拆成高频的小批次。即使单条消息处理需要 100 毫秒,50 条也只需要 5 秒钟。主线程处理完能快速回到下一次 poll() 循环,避免超出 5 分钟超时。

2. 调大间隔参数

根据最坏的业务场景,比如下游服务宕机、网络抖动等,合理调大最大 poll 间隔时间:

max.poll.interval.ms: 600000

不过这通常要配合 max.poll.records 一起微调,不建议设得无限大。因为如果主线程真的彻底死锁卡死了,服务端需要等 10 分钟才能发现,这期间对应的分区就会被其他消费者接管,造成数据消费中断。

3. 异步多线程消费

如果单条业务逻辑确实极其耗时,无论怎么调小 records 都不行,那就不要在 Kafka 的消费主线程里直接干活。

可以只让主线程负责拉数据,拉完立刻丢进自定义的 Java 线程池里去并行处理,让主线程秒级回到下一次 poll()

不过这个方案会引入消息丢失和重复消费的风险:

自动提交可能导致消息丢失:如果开启了自动提交,主线程下一次 poll() 就会把上一批数据对应的 offset 提交。但此时线程池里可能还有一堆任务在排队。一旦服务此时重启或崩溃,这些还没来得及处理的消息就彻底丢失了。

手动提交的乱序与重复消费:要知道线程池里的并发线程是乱序完成的,如果在子线程里直接提交 offset,会发生 offset 覆盖导致重复消费。要解决这个问题,需要自己设计一套滑窗提交机制或者基于每个 Partition 单独配内存队列,开发成本非常高。

说在最后

分布式系统对活着其实有两种定义:一个是进程存活,代表进程在,端口通,心跳还在跳;另一个是业务存活,代表业务还能响应外部输入,推动数据往下走。

写业务消费逻辑时,不能仅凭进程在,日志在刷就认为消费者运行正常,需要合理评估 max.poll.recordsmax.poll.interval.ms 的配置,避免因为单次处理耗时过长导致不必要的 Rebalance。

上次编辑于: