【问题标题】:Understanding reduceByKey function definition Spark Scala理解 reduceByKey 函数定义 Spark Scala
【发布时间】:2017-03-27 02:28:35
【问题描述】:

spark 中Pair RDDreduceByKey 函数具有以下定义:

def reduceByKey(func: (V, V) => V): RDD[(K, V)]

我了解reduceByKey 采用参数函数将其应用于键的值。我想了解的是如何阅读这个定义,其中函数将 2 个值作为输入,即(V, V) => V。不应该是V => V,就像mapValues 函数一样,该函数应用于值V 以产生U,它是相同或不同类型的值:

def mapValues[U](f: (V) ⇒ U): RDD[(K, U)]

这是因为reduceByKey 一次应用于所有值(对于同一个键),而mapValues 一次应用于每个值(与键无关)?在这种情况下应该它被定义为类似(V1, V2) => V

【问题讨论】:

    标签: scala apache-spark


    【解决方案1】:

    ...不应该是V => V,就像mapValues...

    不,它们完全不同。回想一下map 函数中有一个不变量,它们返回一个IterableListArray 等),其length 与原始列表(映射的列表)相同。另一方面,reduce 函数 aggregatecombine 所有元素,在这种情况下,reduceByKey 通过应用函数组合对或值,此定义来自称为monoid 的数学概念。你可以这样看,你通过应用函数组合列表的两个第一个元素,并且该操作的结果应该与第一个元素的类型相同,与第三个元素一起操作,依此类推,直到你最终只有一个元素。

    【讨论】:

      【解决方案2】:

      mapValues 应用 f: (V) ⇒ U 将 RDD 中对的每一第二部分转换为应用 f: (V, V) => V

      val data = Array((1,1),(1,2),(1,4),(3,5),(3,7))
      val rdd = sc.parallelize(data)
      rdd.mapValues(x=>x+1).collect
      // Array((1,2),(1,3),(1,5),(3,6),(3,8))
      rdd.reduceByKey(_+_).collect
      // Array((1,7),(3,12))
      

      【讨论】:

        猜你喜欢
        • 2016-08-26
        • 1970-01-01
        • 2018-08-24
        • 2016-02-25
        • 1970-01-01
        • 2023-03-11
        • 2014-07-19
        • 2018-10-23
        • 2016-09-15
        相关资源
        最近更新 更多