【问题标题】:FLINK: Kafka Source - restart policy when a new topic is discovered at restartFLINK:Kafka Source - 重启时发现新主题时的重启策略
【发布时间】: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 属性是否正在启动?

【问题讨论】:

    标签: apache-kafka apache-flink


    【解决方案1】:

    据我所知,目前FlinkKafkaConsumer 就是这样工作的,如果它从保存点恢复,所有不属于保存点的主题都将自动设置EARLIEST 偏移量。这很可能是一个错误,所以我正在为此创建一个错误报告。

    【讨论】:

    • 您是否介意在打开该错误后分享它?另外:您能指出为任何未知主题设置EARLIEST 的代码吗?
    • 当然可以。这是 jira 的链接:issues.apache.org/jira/browse/FLINK-21219。代码本身在方法open()中的FlinkKafkaConsumerBase中。
    猜你喜欢
    • 2012-04-03
    • 2012-11-06
    • 2016-01-23
    • 1970-01-01
    • 1970-01-01
    • 2018-02-23
    • 1970-01-01
    • 1970-01-01
    • 2017-04-27
    相关资源
    最近更新 更多