尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

深耕网站建设与建站编程的一线实战洞察。

消费线程正常打印日志,为什么 Kafka 却判定宕机、强行 Rebalance?

消费线程正常打印日志,为什么 Kafka 却判定宕机、强行 Rebalance? 线上发版或者业务高峰期有时会遇到一个奇怪的现象消费端的服务进程没挂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.records和max.poll.interval.ms的配置避免因为单次处理耗时过长导致不必要的 Rebalance。
返回列表