【发布时间】: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