【问题标题】:Apache Flink DataStream API doesn't have a mapPartition transformationApache Flink DataStream API 没有 mapPartition 转换
【发布时间】:2016-04-12 14:10:24
【问题描述】:

Spark DStream 有 mapPartition API,而 Flink DataStream API 没有。有没有人可以帮忙解释一下原因。我想做的是在Flink上实现一个类似于SparkreduceByKey的API。

【问题讨论】:

    标签: apache-flink


    【解决方案1】:

    Flink 的流处理模型与以小批量为中心的 Spark Streaming 截然不同。在 Spark Streaming 中,每个 mini batch 都像常规批处理程序一样在有限的数据集上执行,而 Flink DataStream 程序持续处理记录。

    在 Flink 的 DataSet API 中,MapPartitionFunction 有两个参数。输入的迭代器和函数结果的收集器。 Flink DataStream 程序中的MapPartitionFunction 永远不会从第一个函数调用中返回,因为迭代器会遍历无穷无尽的记录流。但是,Flink 的内部流处理模型要求用户函数返回以检查函数状态。因此,DataStream API 不提供mapPartition 转换。

    为了实现类似于 Spark Streaming 的reduceByKey 的功能,您需要在流上定义一个键控窗口。 Windows 将流离散化,这有点类似于小批量,但 Windows 提供了更多的灵活性。由于窗口的大小是有限的,您可以调用reduce该窗口。

    这可能看起来像:

    yourStream.keyBy("myKey") // organize stream by key "myKey"
              .timeWindow(Time.seconds(5)) // build 5 sec tumbling windows
              .reduce(new YourReduceFunction); // apply a reduce function on each window
    

    DataStream documentation 展示了如何定义各种窗口类型并解释了所有可用的功能。

    注意: DataStream API 最近已重新设计。该示例假定最新版本 (0.10-SNAPSHOT) 将在未来几天内发布为 0.10.0。

    【讨论】:

    • 您提供的'reduceByKey'的解决方案似乎类似于'reduceByKey'之外的spark中的'GroupByKey'。 databricks.gitbooks.io/databricks-spark-knowledge-base/content/…
    • 不,Flink 的 reduce() 像 Spark 的 reduceByKey 一样适用于组上的成对归约函数。不过组定义有点不同,因为 Flink 以小批量的方式使用 windows 和 Spark 键值对。在 Flink 中没有直接等效于 Spark 的groupByKey,因为这意味着需要在内存中物化整个组,这可能会导致 OutOfMemoryErrors 并杀死 JVM。 Flink 提供groupReduce() 来消费流式迭代器。
    • 我看到 Flink 的 reduce() 适用于可组合。 Flink DataStream 没有 reduceGroup API 作为 mapPartition 是不是类似的原因?
    • 是和否:-)。是的,对于非窗口流,因为它们的长度是无限的,groupReduce 方法永远不会返回。对于窗口流(窗口长度有限)不适用,其中apply() 方法本质上是一个groupReduce 函数,它还提供窗口的一些元数据。
    • 谢谢,我好像找不到你说的这种申请方法。
    【解决方案2】:

    假设您的输入流是单分区数据(比如字符串)

    val new_number_of_partitions = 4
    
    //below line partitions your data, you can broadcast data to all partitions
    val step1stream = yourStream.rescale.setParallelism(new_number_of_partitions)
    
    //flexibility for mapping
    val step2stream = step1stream.map(new RichMapFunction[String, (String, Int)]{
      // var local_val_to_different_part : Type = null
      var myTaskId : Int = null
    
      //below function is executed once for each mapper function (one mapper per partition)
      override def open(config: Configuration): Unit = {
        myTaskId = getRuntimeContext.getIndexOfThisSubtask
        //do whatever initialization you want to do. read from data sources..
      }
    
      def map(value: String): (String, Int) = {
        (value, myTasKId)
      }
    })
    
    val step3stream = step2stream.keyBy(0).countWindow(new_number_of_partitions).sum(1).print
    //Instead of sum(1), you can use .reduce((x,y)=>(x._1,x._2+y._2))
    //.countWindow will first wait for a certain number of records for perticular key
    // and then apply the function
    

    Flink 流是纯流(不是批处理的)。看看 Iterate API。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2020-09-20
      • 1970-01-01
      • 2021-05-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多