【问题标题】:Passing elements back to the input stream, after processing, in Flink?在 Flink 中处理后将元素传回输入流?
【发布时间】:2020-02-20 13:24:55
【问题描述】:

场景:

我有来自传感器的事件流。事件可以是 T-typeJ-Type

  • T 类事件有事件发生时间戳。
  • J 型事件具有开始和结束时间戳。

根据J-Type事件的开始和结束时间戳,对时间范围内的所有T-type事件应用聚合逻辑,并将结果写入DB。

为此,我创建了一个自定义触发器,它在收到 J-Type 事件时触发。在我的自定义 ProcessWindowFunction 中,我正在执行聚合逻辑和时间检查。

但是,可能存在一种情况,即 T 型事件不在当前 J 型事件的时间范围内。 在这种情况下,T 类事件应该在清除当前窗口之前被推送到下一个窗口。

解决方案的想法:

  1. 在自定义窗口处理函数中将未处理的 T 型事件推送到 Kinesis 流(源)中。 (最坏情况的解决方案)

  2. 使用 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


【解决方案1】:

这里有两个解决方案。它们的基本行为或多或少是相同的,但您可能会发现其中一个更容易理解、维护或测试。

至于您的问题,不,没有办法循环回(重新排队)未使用的事件而不将它们推回 Kinesis。但是只要坚持到需要它们就可以了。

解决方案 1:使用 RichFlatMapFunction

当 T 型事件到达时,将它们附加到 ListState 对象。当 J 型事件到达时,将列表中所有匹配的 T 型事件收集到输出,并更新列表以仅保留那些将属于以后的 J 型事件的 T 型事件。

解决方案 2:将 GlobalWindows 与自定义 Trigger 和 Evictor 一起使用

除了您已经完成的工作之外,实现一个Evictor,它(在窗口被 FIREd 之后)从窗口中只删除 J 类型事件和所有匹配的 T 类型事件。

更新:清除旧密钥/失效传感器的状态

使用解决方案 1,您可以使用 state TTL 安排清除与死键相关的任何非活动状态。或者您可以使用KeyedProcessFunction 而不是RichFlatMapFunction,并使用计时器来完成同样的事情。

使用窗口 API 管理陈旧密钥的状态可能不那么简单,但对于解决方案 2,我相信您可以扩展自定义触发器以包含将清除窗口的超时。如果您在 ProcessWindowFunction 中使用了全局状态,则需要依靠状态 TTL 来清理它。

【讨论】:

  • 嗨,大卫,感谢您的回复。我对 window 的理解是他们是 Akka 演员。如果传感器失效,sensorId 就会过时。因此,如果我们没有针对 windows 的清除策略,这些过时的 actor 将消耗主内存。
  • 重新排队是我们案例中唯一的最佳解决方案,因为它简化了流程并减少了极端情况。如果有任何方法可以将文件直接发送回输入流(无需 kinesis),请分享该方法。
  • 我还想再有一个侧流(kinesis 或任何其他建议),它将接收来自 evictor() 的这些无效元素,并将与主输入流进行 union()。请分享您对此的看法?
  • 我非常喜欢解决方案 1(有状态平面图或流程函数);这就是我会做的。
  • 我看不出重新排队如何解决传感器失效的问题。您最终可能会看到最后几个 T 类事件无休止地循环回到 Kinesis。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2019-11-26
  • 1970-01-01
  • 2017-02-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多