【发布时间】:2020-01-17 00:20:42
【问题描述】:
我有两个 flink dataStream。例如:dataStream1 和 dataStream2。我想将两个流合并为 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