【问题标题】:Consume from two flink dataStream based on priority or round robin way基于优先级或轮询方式从两个 flink dataStream 中消费
【发布时间】:2020-01-17 00:20:42
【问题描述】:

我有两个 flink dataStream。例如:dataStream1dataStream2。我想将两个流合并为 1 个流,以便我可以使用相同的处理函数来处理它们,因为 dataStream 的 dag 是相同的。

截至目前,我需要对任一流的消息消费具有同等优先级。 dataStream2 的生产者每分钟产生 10 条消息,而 dataStream1 的生产者每秒产生 1000 条消息。此外,dataStreams.DataSteam2 的 dataTypes 是相同的,更多的是应该尽快使用的高优先级队列。 dataStream1和dataStream2的消息没有关系

dataStream1.union(dataStream2) 是否会生成一个包含两个 Streams 元素的 Stream?

【问题讨论】:

  • 欢迎您!究竟是什么问题?
  • 数据流从何而来?直接来自源组件?
  • dataStreams 是 pulsar topic 的源组件。
  • @Christophe Does .union() 将产生流,这将是两个数据流的循环。
  • @NischalKumar union() 没有引入任何法规 IIRC。因此,如果您的一个来源比另一个更快地产生元素,那么它就不会调节流量。

标签: apache-flink flink-streaming


【解决方案1】:

可能是这个问题最简单的解决方案,但也不是最有效的解决方案,具体取决于您的数据源的确切规范,它可能是将两个流连接起来。在此解决方案中,您可以使用CoProcessFunction,它将为每个连接的流调用单独的方法。

在此解决方案中,您可以简单地缓冲一个流的元素,直到可以生成它们(例如以循环方式)。但请记住,如果源产生事件的频率之间存在很大差异,这可能会非常低效。

【讨论】:

  • 当每个连接的流都需要单独的业务逻辑时,CoProcessFunction 才有意义。就我而言,业务逻辑是完全一样的。
  • 为什么不让CoProcessFunction中的每个方法调用一个通用的方法来处理记录呢?你在这里的评论让我觉得你的问题比上面描述的要多。
  • 当数据类型不同时,使用 CoProcessFunction 是有意义的。在我的用例中,dataSram1 与 dataStream2 具有相同的数据类型,但 dataStream1 生产者每秒产生 10^3 条消息,而 dataStream2 每分钟产生 10 条消息。我想要一个循环策略,以便 dataSteam2 也与 dataSteam1 消息一起使用。由于dataStream2的速率非常低,它应该在它到达的那一刻被处理。
  • @NischalKumar 不完全是,因为 CoProcessFunction 实际上是由于 Collector 的存在而允许您不为给定流发出输出而只是缓冲它并等待来自另一个流的输入。
【解决方案2】:

听起来这两个DataStreams 具有不同类型的元素,尽管您没有明确指定。如果是这种情况,则通过MapFunction 在每个流上创建Either<stream1 type, stream2 type>,然后在两个流上创建union()。你不会得到两者的精确混合,因为 Flink 会交替使用每个流的网络缓冲区。

如果您真的想要很好地混合流,那么(正如其他人所指出的)您需要通过状态缓冲传入的元素,并应用一些启发式方法来避免由于任何原因(例如不同的网络延迟,或两个源之间的性能更可能不同)您在两个流之间具有非常不同的数据速率。

【讨论】:

  • 两个数据流的数据类型相同
【解决方案3】:

您可能希望使用实现InputSelectable 接口的自定义运算符,以减少所需的缓冲量。我在下面提供了一个示例,该示例在没有任何缓冲的情况下实现了交错,但请务必阅读docs 中的警告,它解释了

...操作员可能会收到一些当前不想处理的数据...

换句话说,不能依赖这个简单的示例来真正按原样工作。

public class Alternate {
    public static void main(String[] args) throws Exception {

        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);

        DataStream<Long> positive = env.generateSequence(1L, 100L);
        DataStream<Long> negative = env.generateSequence(-100L, -1L);

        AlternatingTwoInputStreamOperator op = new AlternatingTwoInputStreamOperator();

        positive
            .connect(negative)
            .transform("Hack that needs buffering", Types.LONG, op)
            .print();

        env.execute();
    }
}

class AlternatingTwoInputStreamOperator extends AbstractStreamOperator<Long>
        implements TwoInputStreamOperator<Long, Long, Long>, InputSelectable {

    private InputSelection nextSelection = InputSelection.FIRST;

    @Override
    public void processElement1(StreamRecord<Long> element) throws Exception {
        output.collect(element);
        nextSelection = InputSelection.SECOND;
    }

    @Override
    public void processElement2(StreamRecord<Long> element) throws Exception {
        output.collect(element);
        nextSelection = InputSelection.FIRST;
    }

    @Override
    public InputSelection nextSelection() {
        return this.nextSelection;
    }
}

还要注意InputSelectable 是在 Flink 1.9.0 中添加的。

【讨论】:

    猜你喜欢
    • 2021-07-23
    • 1970-01-01
    • 1970-01-01
    • 2013-09-20
    • 2021-01-09
    • 1970-01-01
    • 1970-01-01
    • 2021-02-11
    • 1970-01-01
    相关资源
    最近更新 更多