【问题标题】:How does Flink handle watermarks with Union operators?Flink 如何处理 Union 算子的水印?
【发布时间】:2020-02-25 21:04:50
【问题描述】:

我从四个 Kinesis 流中读取数据。每个流中的数据是不同的数据类型。在读入所有四个流之后,我分配时间戳和水印,并聚合来自每个流的数据。四个聚合的结果都使用相同的通用对象输出。我想合并来自四个流的结果,以便可以将合并的流发送到一个 ProcessFunction。这基本上允许我像 CoProcessFunction 一样使用 ProcessFunction,但我将能够处理来自两个以上流的数据(在这种情况下,ProcessFunction 将接收来自所有四个单独流的聚合)。

但是,我担心这可能无法很好地处理水印。如果一个流需要更长的时间来处理或以某种方式落后,如果所有水印在联合中向前传递并且其中一个流领先于其他流,则它的聚合可能无法进入处理函数。如果是这种情况,那么处理函数的水印将是它从四个单独的流中看到的水印的最大值。

我的问题是:联合运营商如何处理水印,联合运营商下游的运营商如何处理这些水印?

另外:如果泛型对象的联合由于水印问题而不起作用,当 Flink 仅支持两个流的 CoProcessFunction 时,将四个不同聚合的结果组合起来的最佳方法是什么?

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    将两个以上的流连接在一起的另一种方法是构建一棵树,该树进行成对连接,直到所有流都连接在一起。要么作为平衡树,像这样:

    A--->
         A+B---->
    B--->
    
                A+B+C+D------------>
    
    C--->
         C+D---->
    D--->
    

    或一次添加一个流,如下所示:

    a--->
         a+b--->
    b--->
                a+b+c--->
         c----->
                         a+b+c+d--->
                d------->
    

    FWIW,FLIP-92 是向 Flink 添加 n 元流运算符的提议,但即使实施,它也可能不会对用户可见,至少一开始是这样。

    【讨论】:

    • 我是否正确地认为我可以使用同一个通用 Java 对象发出四个单独的聚合,这样我就可以将它们合并并在一个 ProcessFunction 中一起处理它们,本质上是创建我自己的 n 元流运算符?您是否发现这种方法存在任何问题,尤其是在四个流之一落后的情况下?
    • 您向工会提出的建议应该可以正常工作。我以前看过它。水印只会随着输入流水印的最小值而进步。
    • 太好了,再次感谢您!总是快速响应并且非常有见地!
    【解决方案2】:

    联合水印的工作原理与并行流的水印一样。这意味着水印始终是来自所有输入流的水印的min。同样代表下游操作符,它们的水印将是所有输入流的min

    老实说,我认为联合并不依赖于水印。但是,如果您出于任何原因想要使用 CoProcessFunction,我可以提供这种有点 hacky 的方式。您可以创建一个 Seq 您已生成的流,然后:

    //Streams defined
    val seq = Seq(stream, stream2, stream3, stream4)
    seq.reduce((stream1, stream2) => stream1.connect(stream2).process(...))
    

    【讨论】:

    • 联合本身不依赖于水印,但联合之后的处理函数将依赖于水印。如果一个流落后于其他流,那么遵循四个流联合的流程函数将如何表现?它会等到滞后的流赶上来才触发计时器吗?文档特别指出 CoProcessFunctions 旨在获取两个流中的最小值,但没有指定一个联合后跟一个 ProcessFunction 对 N 个流的作用。
    • 嗯,确实如此 :) 这里有一个关于并行流中的水印的部分:ci.apache.org/projects/flink/flink-docs-stable/dev/…。这意味着,如果其中一个流滞后,那么它可以有效地防止工作时间的进展。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-07-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多