【发布时间】:2020-02-25 21:04:50
【问题描述】:
我从四个 Kinesis 流中读取数据。每个流中的数据是不同的数据类型。在读入所有四个流之后,我分配时间戳和水印,并聚合来自每个流的数据。四个聚合的结果都使用相同的通用对象输出。我想合并来自四个流的结果,以便可以将合并的流发送到一个 ProcessFunction。这基本上允许我像 CoProcessFunction 一样使用 ProcessFunction,但我将能够处理来自两个以上流的数据(在这种情况下,ProcessFunction 将接收来自所有四个单独流的聚合)。
但是,我担心这可能无法很好地处理水印。如果一个流需要更长的时间来处理或以某种方式落后,如果所有水印在联合中向前传递并且其中一个流领先于其他流,则它的聚合可能无法进入处理函数。如果是这种情况,那么处理函数的水印将是它从四个单独的流中看到的水印的最大值。
我的问题是:联合运营商如何处理水印,联合运营商下游的运营商如何处理这些水印?
另外:如果泛型对象的联合由于水印问题而不起作用,当 Flink 仅支持两个流的 CoProcessFunction 时,将四个不同聚合的结果组合起来的最佳方法是什么?
【问题讨论】:
标签: apache-flink flink-streaming