【问题标题】:Flink: What happens when ProcessAllWindowFunction takes more time than the TumblingProcessingTimeWindows defined in windowAll()Flink:当 ProcessAllWindowFunction 花费的时间超过 windowAll() 中定义的 TumblingProcessingTimeWindows 时会发生什么
【发布时间】:2021-10-01 08:41:05
【问题描述】:

我有一个 ProcessAllWindowFunction 实现(请参阅下面代码中的 AttributeBackLogEvents()),它有很多 I/O,可能需要 30 多秒。 windowAll() 是使用 30 秒的 TumblingProcessingTimeWindows 窗口化数据。

attributedStream
            .windowAll(TumblingProcessingTimeWindows.of(Time.seconds(30)))
            .process(new AttributeBackLogEvents())
            .forceNonParallel()
            .addSink(ConfluentKafkaSink.createKafkaSinkFromApplicationProperties())
            .name("Enriched Event kafka topic sink");

AttributeBackLogEvents 根据传递的迭代从 MySQL 获取一组事件,并在一些处理后删除 MySQL 的一些获取事件。我看到当前窗口正在获取的记录(理想情况下应该在下一个窗口触发之前删除)也被下一个窗口获取,这意味着即使当前窗口正在处理下一个窗口也会触发。

我的问题是:

  1. 是否有可能 AttributeBackLogEvents 仍在运行并触发下一个窗口?
  2. 如果是这样,那么在当前窗口处理完成之前,我如何强制执行此操作,下一个窗口不应触发。

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    这个 Q 并没有描述逻辑中发生的事情,而是在概念上: 您的窗口意味着“源数据的时间范围”,因此在下一个窗口开始之前,任何处理都无法完全完成。

    可能有一种方法,但对于流式处理工具,像 MySQL 这样的源通常被视为参考数据(您通常希望经常阅读),除非您正在执行更改数据捕获。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2018-10-25
      • 1970-01-01
      • 1970-01-01
      • 2023-01-30
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2016-02-20
      相关资源
      最近更新 更多