【发布时间】:2021-05-03 11:01:21
【问题描述】:
我有一个 flink 作业,通过 KafkaSource 配置为监听主题的正则表达式,例如:
val topicPattern = "^(topic1|topic2|topic3)$"
Kafka 消费者开始位置配置设置为 startFromLatest,如下所示:
val myConsumer = new FlinkKafkaConsumer<>(topicPattern, someProperties);
myConsumer.setStartFromLatest();
我们通过配置传递 topicPattern,有时会发生一个新的 kafka 生产者生成数据,比如说topic4,然后我们更新配置添加这个新主题并使用保存点重新启动作业。
在这种情况下,我们注意到 kafka 源从头开始读取这个新主题。有没有人能解释为什么? Kafka auto.offset.reset 属性是否正在启动?
【问题讨论】: