Kafka Streams再平衡期间状态存储重建导致分区处理停滞问题咨询
问题1解答
在Kafka Streams 2.8版本之前默认使用的Eager再平衡协议下,不仅part1会停滞,所有分区的处理都会在再平衡触发后全部暂停,直到再平衡完成、所有实例完成对应分区的状态恢复才会恢复处理。
2.8版本之后默认启用的增量协作再平衡协议下,仅涉及迁移的part1会在所有权移交后暂停处理,直到instance3完成状态恢复才会继续消费,这段时间part1的消息会正常堆积在Kafka集群不会丢失,只是处理延迟升高。不管哪种协议,part1的处理确实会在状态恢复阶段停滞,因为Kafka Streams要保证状态一致性,必须将本地状态恢复到对应偏移量位置才会启动新消息的处理。
问题2解答
可以通过以下几个维度估算恢复时间:
- 先统计
part1对应changelog分区的有效数据量:可以用kafka-log-dirs.sh工具查看分区的实际存储大小,如果changelog开启了日志压缩,需要排除被标记为删除的无效条目,只统计有效状态的总大小。 - 参考历史状态恢复速率:Kafka Streams自带
state-restoration-rate(状态恢复速率)、remaining-restoration-time(剩余恢复时间预估)的监控指标,直接用有效数据量除以历史平均恢复速率就能得到大致的耗时,比如历史恢复速率为5MB/s,10GB有效数据的话预估耗时约34分钟,可额外预留20%~30%的冗余量。 - 考虑底层资源限制:如果用RocksDB作为状态存储,需要结合本地磁盘的写入IO性能、实例到Kafka集群的网络带宽上限来估算;如果开启了状态Standby副本,恢复时会先从Standby实例拷贝全量状态文件,再追少量最新的changelog数据,耗时会比全量读changelog少很多,此时可以按文件拷贝速率加增量changelog读取时间估算。
问题3解答
和再平衡协议有关:
- 如果用的是2.8版本之前默认的Eager再平衡协议,所有实例都会在再平衡触发时停止处理、交出所有已分配的分区,直到整个再平衡流程全量完成才会恢复,所以未涉及迁移的
part2、part3也会停止处理。 - 如果用的是2.8及以后默认的增量协作再平衡协议,仅涉及迁移的分区会暂停处理,
part2、part3的分配关系没有变化,对应的实例可以无停机持续处理数据,完全不受这次扩容的影响。
内容的提问来源于stack exchange,提问作者Nementaarion
相关产品推荐
相关产品推荐

