修复Bug重启Flink应用后,是否会基于Kafka偏移量进行消费?
问题解答
重启后的Flink应用会从Kafka Broker上已提交的偏移量开始消费,不会从头读取Kafka队列,原因如下:
- 当
commit.offsets.on.checkpoint配置为true时,Flink每次成功完成Checkpoint后,都会将当前消费的Kafka偏移量提交到Kafka Broker的__consumer_offsets主题中,这个偏移量的存储完全独立于Flink自身的Checkpoint或Savepoint。 - 此次重启丢失的是Flink内部的Checkpoint数据,但Kafka端保存的已提交偏移量不受任何影响。
- 新部署的Flink应用启动时,会默认读取Kafka Broker中对应消费者组的已提交偏移量作为消费起点;只有当Kafka中没有该消费者组的提交偏移量,且
auto.offset.reset配置为earliest时,才会从头消费——显然你的场景不满足这个条件。
内容的提问来源于stack exchange,提问作者Dhruv Chadha
相关产品推荐
相关产品推荐

