【问题标题】:How to specify two sources, one process operator and one sink operator in flink applicationflink应用中如何指定两个source,一个process operator,一个sink operator
【发布时间】: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


    【解决方案1】:

    API 在您如何设置处理管道方面提供了很大的灵活性。如果您想将相同的逻辑应用于多个来源,您可以这样做:

    env.addSource(new MySource1())
      .process(new MyProcess())
      .addSink(new MySink())
    
    env.addSource(new MySource2())
      .process(new MyProcess())
      .addSink(new MySink())
    
    env.execute()
    

    或者如果这样做更有意义,您可以合并两个流,然后处理组合流(或这些方法的某种组合):

    stream1.union(stream2)
      .process(...)
      .addSink(...)
    

    如果您想分叉流并对每个副本应用不同的操作,也可以反过来做:

    val stream: DataStream[T] = env.addSource(new MySource())
    
    stream.process(new MyProcess1())
      .addSink(new MySink1())
    
    stream.process(new MyProcess2())
      .addSink(new MySink2())
    
    env.execute()
    

    哇,Flink 1.3 已经三年多了!

    【讨论】:

    • 感谢 David @david-anderson 的完美回答!是的,1.3 太旧了,>
    猜你喜欢
    • 1970-01-01
    • 2010-12-14
    • 1970-01-01
    • 1970-01-01
    • 2012-11-28
    • 2010-10-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多