Kafka Streams 3.7.0 GlobalKTable checkpoint生成异常问题咨询
Kafka Streams 2.5.0 vs 3.7.0 GlobalKTable Checkpoint机制退化问题
我们对比了Kafka Streams 2.5.0与3.7.0版本中GlobalKTable的checkpoint机制,发现3.7.0存在功能退化问题:首次消费完主题数据后,若没有额外处理10000条记录,.checkpoint文件不会生成。
2.5.0版本行为
- 流处理任意记录后,最多在
max.commit.interval.ms(默认30秒)内生成.checkpoint文件 - 消费完主题后,重启应用可直接通过该文件恢复偏移量,无需重新消费全量数据
3.7.0版本退化表现
- 代码仅在
totalOffsetDelta > OFFSET_DELTA_THRESHOLD_FOR_CHECKPOINT(即10_000)时才强制生成checkpoint totalOffsetDelta仅在记录添加到存储时递增,首次消费主题的过程中该值无法超过阈值,因此不会生成.checkpoint文件- 即使一次性消费6000万条记录,仍需额外处理10000条记录才能触发生成checkpoint
业务场景影响
在主题消费完成后不再产生新记录的业务场景中,.checkpoint文件永远无法生成。重启Kafka Streams应用时,因无checkpoint可恢复,需重新消费全部6000万条记录,耗时约4分钟才能恢复到最新状态。
内容的提问来源于stack exchange,提问作者paul
相关产品推荐
相关产品推荐

