【发布时间】:2021-07-28 05:40:03
【问题描述】:
我对 Flink 和 Kafka 还很陌生,并且有一些用 Scala 编写的数据聚合作业,这些作业在 Apache Flink 中运行,这些作业使用来自 Kafka 的数据执行聚合并将结果返回给 Kafka。
我需要作业来使用作业运行时创建的与模式匹配的任何新 Kafka 主题的数据。我通过为我的消费者设置以下属性来完成这项工作
val properties = new Properties()
properties.setProperty(“bootstrap.servers”, “my-kafka-server”)
properties.setProperty(“group.id”, “my-group-id”)
properties.setProperty(“zookeeper.connect”, “my-zookeeper-server”)
properties.setProperty(“security.protocol”, “PLAINTEXT”)
properties.setProperty(“flink.partition-discovery.interval-millis”, “500”);
properties.setProperty(“enable.auto.commit”, “true”);
properties.setProperty(“auto.offset.reset”, “earliest”);
val consumer = new FlinkKafkaConsumer011[String](Pattern.compile(“my-topic-start-.*”), new SimpleStringSchema(), properties)
消费者工作正常并使用以“my-topic-start-”开头的现有主题的数据
当我第一次针对新主题发布数据时,例如“my-topic-start-test1”,我的消费者直到主题创建后 500 毫秒后才识别该主题,这是基于特性。 当消费者识别出主题时,它不会读取发布的第一条数据记录,而是开始有效地读取后续记录,因此每次针对新主题发布数据时,我都会丢失第一条数据记录。
是否有我遗漏的设置或者 Kafka 的工作方式?任何帮助将不胜感激。
谢谢 沙文
【问题讨论】:
标签: scala apache-kafka apache-flink kafka-consumer-api