【发布时间】:2020-06-09 18:24:37
【问题描述】:
我正在尝试使用 CoProcessFunction 加入 2 个流。输入流A之一是实时生成的。然而,另一个输入流 B 由一个延迟 1 天的每日计划作业加载,这意味着今天放入流中的事件始终具有昨天的事件时间。
话虽如此,流 B 的水印总是比 A 的水印晚约 1 天,所以我想来自 A 的很多事件将被缓冲在内存中。我想知道是否有办法解决问题。一些额外的背景,流 A 和 B 都是 kinesis 流(我使用的是FlinkKinesisConsumer),保留期 = 7 天。
提前致谢!
【问题讨论】:
-
你能告诉更多/也许提供一个连接看起来如何的例子吗?你想加入来自 B 的时间戳相等/更低的时间戳吗?
-
当然,举个简单的例子:两个流中的事件都有一个字段
dataId,所以对于流B(延迟的那个)中的给定事件X,我想找到0个或多个流 A 中具有相同dataId且事件时间不晚于事件X的事件时间后 5 分钟的事件。 -
对机制的期望是什么?你不想要任何缓冲?
-
我知道未处理的事件将被缓冲在内存中,在我的情况下,来自当天生成的流 A 的事件将被缓冲。我只是担心应用程序可能会耗尽内存并崩溃。
标签: apache-flink flink-streaming