YuF*_*hen 3 apache-flink flink-streaming
我想让 Windows 在滚动处理时间达到 100 或每 5 秒后完成?也就是说当元素达到100时,触发Windows计算,但如果元素没有达到100,但时间过去了5秒,也会触发Windows计算,就像下面两个触发器的组合:
.countWindow(100)
.window(TumblingProcessingTimeWindows.of(Time.seconds(5)))
使用当前的 Flink API 没有超级简单的方法来做到这一点。
您的用例需要状态(用于计数)和计时器的组合。您可以使用自定义Trigger或使用ProcessFunction来完成此操作。
对于使用 windows 和自定义触发器的方法,查看ProcessingTimeTrigger 和 CountTrigger的实现会有所帮助,因为您基本上希望将两者混合。
ProcessFunction 是一个较低级别的构建块,它结合了托管状态和计时器,这正是您所需要的,所以这可能更容易,特别是如果您已经知道如何使用Flink 的托管状态。
顺便说一句,在线 Flink 培训包括学习如何使用 ProcessFunctions 的材料。
| 归档时间: |
|
| 查看次数: |
1889 次 |
| 最近记录: |