我们可以在 Flink 中结合计数和处理时间触发器吗?

YuF*_*hen 3 apache-flink flink-streaming

我想让 Windows 在滚动处理时间达到 100 或每 5 秒后完成?也就是说当元素达到100时,触发Windows计算,​​但如果元素没有达到100,但时间过去了5秒,也会触发Windows计算,​​就像下面两个触发器的组合:

.countWindow(100)

.window(TumblingProcessingTimeWindows.of(Time.seconds(5)))

Dav*_*son 5

使用当前的 Flink API 没有超级简单的方法来做到这一点。

您的用例需要状态(用于计数)和计时器的组合。您可以使用自定义Trigger或使用ProcessFunction来完成此操作。

对于使用 windows 和自定义触发器的方法,查看ProcessingTimeTrigger 和 CountTrigger的实现会有所帮助,因为您基本上希望将两者混合。

ProcessFunction 是一个较低级别的构建块,它结合了托管状态和计时器,这正是您所需要的,所以这可能更容易,特别是如果您已经知道如何使用Flink 的托管状态。

顺便说一句,在线 Flink 培训包括学习如何使用 ProcessFunctions 的材料。