【问题标题】:How to connect more than 2 streams in Flink?如何在 Flink 中连接 2 个以上的流?
【发布时间】:2020-10-30 21:28:30
【问题描述】:

我有 3 个不同类型的键控数据流。

DataStream<A> first;
DataStream<B> second;
DataStream<C> third;

每个流都有自己定义的处理逻辑,并在它们之间共享一个状态。只要数据在任何流中可用,我想连接这 3 个流以触发各自的处理功能。可以连接两个流。

first.connect(second).process(<CoProcessFunction>)

我不能使用联合(允许多个数据流),因为类型不同。我想避免创建包装器并将所有流转换为相同的类型。

【问题讨论】:

    标签: apache-flink flink-streaming


    【解决方案1】:

    除了联合之外,标准方法是在级联中使用连接,例如,

    first.connect(second).process(...).connect(third).process(...)

    您将无法在一个地方在所有三个流之间共享状态。您可以让第一个流程函数输出后续流程函数需要的任何内容,但第三个流将无法影响第一个流程函数中的状态,这对于某些用例来说是个问题。

    另一种可能性可能是利用较低级别的机制——请参阅FLIP-92: Add N-Ary Stream Operator in Flink。但是,此机制是供内部使用的(Table/SQL API 将其用于 n 路连接),因此需要谨慎对待。有关详细信息,请参阅mailing list discussion。我提到这一点是为了完整性,但我怀疑在进一步开发界面之前这是一个好主意。

    您可能还想查看stateful functions api,它克服了数据流 api 的许多限制。

    【讨论】:

      【解决方案2】:

      包装方法还不错,真的。您可以创建一个类似于 Flink 现有 Either&lt;Left, Right&gt;EitherOfThree&lt;T1, T2, T3&gt; 包装类,然后在单个函数中处理这些记录的流。比如:

          DataStream <EitherOfThree<A,B,C>> combo = first.map(r -> new EitherOfThree<A,B,C>(r, null, null))
              .union(second.map(r -> new EitherOfThree<A,B,C>(null, r, null)))
              .union(third.map(r -> new EitherOfThree<A,B,C>(null, null, r)));
          combo.process(new MyProcessFunction());
      

      Flink 的 Either 类有一个更优雅的实现,但是对于您的用例来说,一些简单的东西应该可以工作。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多