【发布时间】:2020-08-01 16:55:29
【问题描述】:
我使用的是 flink 1.3,我定义了两个流源,它们会发出相同的事件以供后续操作符(我定义的流程操作符和接收器操作符)处理
但看起来在source-process-pink管道中,我只能指定一个源,我会问如何指定两个或多个源并执行相同的进程和接收器
object FlinkApplication {
def main(args: Array[String]): Unit = {
val env = StreamExecutionEnvironment.getExecutionEnvironment
env.setStreamTimeCharacteristic(TimeCharacteristic.ProcessingTime)
env.addSource(new MySource1()) //How to MySource2 here?
.setParallelism(1)
.name("source1")
.process(new MyProcess())
.setParallelism(4)
.addSink(new MySink())
.setParallelism(2)
env.execute("FlinkApplication")
}
}
【问题讨论】:
标签: apache-flink