【问题标题】: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 提供的用于对有界流进行操作的机制是将setBoundedsetUnbounded 与构建器一起使用,如

    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())
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2019-02-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-01-17
      • 1970-01-01
      • 2015-10-12
      • 2017-02-13
      相关资源
      最近更新 更多