【问题标题】:What is the correct way of connecting a ProcessWindowFunction with a broadcast stream in Flink?在 Flink 中将 ProcessWindowFunction 与广播流连接的正确方法是什么?
【发布时间】:2020-04-16 11:47:41
【问题描述】:

我有一个运行多个模型的 flink 管道,所以窗口看起来像这样:

DataStream<WindowDeviationResult> aggregatedWindow = keyedStream
                                                        .timeWindow(Time.seconds(window_duration))
                                                        .aggregate( model.getWindowAgreggator(), 
                                                                    model.getWindowProcessor());

我需要将来自另一个流的状态发送到 ProcessWindowFunction 运算符(最后一个)。通常,我会在之前进行连接,然后实现 proceessElementprocessBroadcastElement。但是因为我将 WindowProcessFuction 传递给 .aggregate 作为第二个参数,所以我不能这样做。您在这里看到哪些选项?

【问题讨论】:

    标签: java stream apache-flink flink-streaming


    【解决方案1】:

    Flink 不支持将广播流连接到窗口操作符。我应该建议使用 KeyedBroadcastProcessFunction 而不是窗口,并实现自己的窗口。通常这并不是特别困难。请参阅 https://stackoverflow.com/a/59823254/2000823 了解可能有助于您入门的示例。

    【讨论】:

    • 他这是一个非常有用的答案,。我唯一担心的是它是否能够处理迟到?
    • 是的,没问题。您必须决定在清除每个窗口的状态之前保持多长时间(即允许的延迟),然后当延迟事件到达时,如果它们的窗口不再存在,您可以发送延迟事件到侧面输出——如果窗口确实存在,你可以处理迟到的事件。
    • 优秀。我最终做了一个不同的解决方案,但是您的评论和解决方案将对以后的问题非常有用!
    猜你喜欢
    • 2018-11-13
    • 2020-08-03
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2014-05-08
    • 2014-07-04
    • 2012-01-05
    • 1970-01-01
    相关资源
    最近更新 更多