【发布时间】: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 的一些获取事件。我看到当前窗口正在获取的记录(理想情况下应该在下一个窗口触发之前删除)也被下一个窗口获取,这意味着即使当前窗口正在处理下一个窗口也会触发。
我的问题是:
- 是否有可能 AttributeBackLogEvents 仍在运行并触发下一个窗口?
- 如果是这样,那么在当前窗口处理完成之前,我如何强制执行此操作,下一个窗口不应触发。
【问题讨论】:
标签: apache-flink flink-streaming