【发布时间】:2020-02-20 13:24:55
【问题描述】:
场景:
我有来自传感器的事件流。事件可以是 T-type 或 J-Type。
- T 类事件有事件发生时间戳。
- J 型事件具有开始和结束时间戳。
根据J-Type事件的开始和结束时间戳,对时间范围内的所有T-type事件应用聚合逻辑,并将结果写入DB。
为此,我创建了一个自定义触发器,它在收到 J-Type 事件时触发。在我的自定义 ProcessWindowFunction 中,我正在执行聚合逻辑和时间检查。
但是,可能存在一种情况,即 T 型事件不在当前 J 型事件的时间范围内。 在这种情况下,T 类事件应该在清除当前窗口之前被推送到下一个窗口。
解决方案的想法:
在自定义窗口处理函数中将未处理的 T 型事件推送到 Kinesis 流(源)中。 (最坏情况的解决方案)
使用 FIRE 代替 FIRE_AND_PURGE,以在整个运行时维护状态。使用元素迭代器删除已处理的元素。 (不推荐,保持无限窗口)
想知道,是否有任何方法可以将未处理的事件直接推送回输入流(无需运动)。 (重新排队)
或者
有什么方法可以在 keyBy Context 中维护状态,以便我们对这些未处理的数据(之前或)与窗口元素一起执行计算。
【问题讨论】:
-
在讨论解决方案之前,我想更好地理解问题。 T 和 J 事件是否在同一个流中?那个流是键控的——如果是这样,通过什么(可能是sensorId?)?您解释了如何同时打开两个窗口。事情有多乱?可以同时打开三个或更多窗口吗?窗口的 J 事件是否总是先于(或跟随)随之而来的 T 事件,或者是否可以进行任何排序?
-
1.是的,它们在同一个流中。 2. 是的,流是keyBy("sensorId")。 3. 同一个“sensorId”不会同时打开两个窗口。第一个窗口在我们收到第一个 J-Type 后关闭。然后打开第二个窗口。这里 J-Type 启动窗口触发器。在 Window 1 中,我们还可以接收 T-types,它应该属于后续窗口。
-
我们根据 T 型事件中的 event_occured 时间戳将 T 型事件与 J 型事件相关联。如果 event_occured 时间戳落在 J-Type 时间范围内,则属于该 J-Type。如果不在时间范围内,则将这些 T 类事件推送到下一个窗口,以检查它是否在那里有效。
-
处理引擎通常根据有向无环图或 DAG 来考虑流程管道。这是处理可以按特定顺序通过函数的地方,函数可以链接在一起,但处理绝不能回到图中的较早点。从此以后,重新排队不再是一种策略。使用
ListState作为我当前场景的解决方案。参考:blog.scottlogic.com/2018/07/06/…
标签: java apache-flink flink-streaming amazon-kinesis amazon-kinesis-analytics