【问题标题】:Terminate a Flink job when using a Kafka Source使用 Kafka Source 时终止 Flink 作业
【发布时间】:2022-10-04 16:57:20
【问题描述】:
当我的生产者完成将其所有消息流式传输到 Kafka,并且在 Flink 完成处理它们之后,我希望能够终止 Flink 作业,这样它就不会继续运行,并且我也可以知道 Flink 何时完成处理所有的数据。我也不能使用批处理,因为我需要 Flink 与我的 Kafka 流并行运行。
通常,Flink 在 DeserializationSchema 类中使用 isEndOfStream 方法来查看它是否应该提前结束(在方法中返回 true 会自动结束工作)。但是,当使用 Kafka 作为 Flink 的源时,新的 KafkaSource 类已弃用在反序列化器中使用 isEndOfStream 方法,并且不再检查它以查看流是否应该结束。还有其他方法可以提前终止 Flink 作业吗?
【问题讨论】:
标签:
apache-kafka
apache-flink
flink-streaming
【解决方案1】:
KafkaSource 提供的用于对有界流进行操作的机制是将setBounded 或setUnbounded 与构建器一起使用,如
KafkaSource<String> source = KafkaSource
.<String>builder()
.setBootstrapServers(...)
.setGroupId(...)
.setTopics(...)
.setDeserializer(...) // or setValueOnlyDeserializer
.setStartingOffsets(...)
.setBounded(...) // or setUnbounded
.build();
setBounded 表示一旦源消耗了通过指定偏移量的所有数据,就应该停止源。
setUnbounded 可用于指示虽然源不应读取超过指定偏移量的任何数据,但它应保持运行。如果在 STREAMING 模式下运行,这允许源参与检查点。
如果你预先知道你想读多少,这很好用。我使用了带有特定时间戳的setBounded,例如,
.setBounded(
OffsetsInitializer.timestamp(
Instant.parse("2021-10-31T23:59:59.999Z").toEpochMilli()))
也像这样
.setBounded(OffsetsInitializer.latest())