【问题标题】:Generating "Heartbeat" Type Events to Push Event Time Forward生成“心跳”类型的事件来推动事件时间
【发布时间】:2023-04-01 23:49:01
【问题描述】:

我正在构建一个 Flink 流应用程序,并且更喜欢使用事件时间,因为它可以确保所有设置的计时器都会在历史数据失败或重放的情况下确定性地触发。事件时间的问题在于,只有在事件发生时时间才会向前移动。我们的数据源(物理传感器)有时会产生非常少的数据,因此有时单个数据点可能会打开五分钟的聚合窗口,但下一个数据点在 20 分钟后,因此窗口关闭并很晚才发出输出记录。

我们提出的解决方案是使用 AWS lambda 函数,该函数计划每 X 分钟运行一次,将虚拟事件输出到 Flink 读取的 Kinesis 流中,从而强制生成水印以提前时间。

我担心的是,这仅在 watermarks 是真正全局的情况下才有效,这意味着 SINGLE 心跳消息可能会导致创建 watermark,从而提前使用源自的数据的 Flink 应用程序中每个操作员/任务的事件时间从这个流。文档让我相信 Flink 并行化从源读取,其中每个并行读取操作符生成自己的水印,然后下游操作符(例如窗口)获取它所看到的各种水印中的最小值。如果是这种情况,这对我来说似乎有问题,因为每个并行水印生成器都需要一个虚拟心跳事件,但我无法控制哪些节点从流中读取我的心跳消息。

所以,我的问题是,下游操作员究竟如何使用水印来提前事件时间,是否可以将单个虚拟消息添加到 kinesis 流以提前整个 Flink 应用程序的事件时间?

如果没有,我该如何强制事件时间向前推进?

【问题讨论】:

  • 更具体地说,我的应用程序从单个流中读取。从该流中,它获取事件,反序列化它们,分配时间戳并生成水印,通过设备 ID 对它们进行键控,然后将它们发送到 ProcessFunction()。我希望能够发送单个心跳消息,该消息创建一个水印,为每个 KeyedProcessFunction 提前事件时间,以确保我在 KeyedProcessFunction 中注册的计时器不会延迟太多处理时间。此心跳消息本质上将具有 event_timestamp,但其 device_id 字段将为空。

标签: apache-flink flink-streaming


【解决方案1】:

你是对的;这里有一个问题。 BoundedOutOfOrdernessTimestampExtractor 实现的标准周期性水印生成器依赖于查看具有更大时间戳的新事件以推进水印。

有几种方法可以解决这个问题:

  1. 在以 1 的并行度运行的任务中运行源和水印分配器(如果需要,然后增加管道其余部分的并行度)。这样一条心跳消息就足够了。

  2. 广播心跳消息。这样每个并行实例都会收到它们,并且它们都可以推进它们的水印。

  3. 实现一个水印生成器,而不是心跳消息,它使用处理时间计时器来人为地推进水印,尽管没有传入事件。示例见https://github.com/aljoscha/flink/blob/6e4419e550caa0e5b162bc0d2ccc43f6b0b3860f/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/functions/timestamps/ProcessingTimeTrailingBoundedOutOfOrdernessTimestampExtractor.java

请注意,这第三种方法不太理想,因为它与处理时间产生了耦合,从而消除了纯事件时间方法的一些核心优势。

如果您使用心跳源,则需要为返回 MAX_WATERMARK 的其他(有时是空闲的)源实现水印生成器。否则,此流中的水印将阻止整个水印。

此外,AWS Lambda 感觉有点矫枉过正。您可以实现一个简单的自定义 Flink 源来创建心跳事件。

【讨论】:

  • 如果 Flink 集群宕机一个小时,并且事件在流中备份,会发生什么?如果我们严格使用事件时间和常规的 BoundedOutOfOrdernessTimestampExtractor,那么我的理解是流中的事件都会得到正确处理,即使挂钟已经远远超过了事件时间戳。如果我们从上面选择选项 3,并且 Flink 在中断后重新上线,水印生成器是否会假设没有传入事件并推进水印,从而阻止处理备份的记录?
  • 是的,没错,选项3有问题。
  • @ChrisATX 你能解决你的问题吗?如果是这样,您能分享一下您采用的解决方案吗?
  • 我使用类似于上面选项 3 中的解决方案解决了我的问题。我添加了一些额外的逻辑,在部署后等待 X 时间,然后再声明源空闲,让 Flink 有时间在强制推进水印之前从数据流中读取中断的地方。
猜你喜欢
  • 2021-02-03
  • 2018-03-03
  • 2013-08-21
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-10-26
  • 2018-07-20
  • 1970-01-01
相关资源
最近更新 更多