【问题标题】:How to calculate the mean of each pair in an RDD consisting of (Key, [Value]) pairs in Spark?如何计算由 Spark 中的 (Key, [Value]) 对组成的 RDD 中每一对的平均值?
【发布时间】:2015-09-04 08:51:39
【问题描述】:

我对 Scala 和 Spark 都很陌生,所以如果我完全错误地处理这个问题,请原谅我。取一个csv文件,过滤,映射后;我有一个 RDD,它是一堆 (String, Double) 对。

(b2aff711,-0.00510)
(ae095138,0.20321)
(etc.)

当我在 RDD 上使用 .groupByKey() 时,

val grouped = rdd1.groupByKey()

得到一个包含一堆 (String, [Double]) 对的 RDD。 (我不知道 CompactBuffer 是什么意思,也许会导致我的问题?)

(32540b03,CompactBuffer(-0.00699, 0.256023))
(a93dec11,CompactBuffer(0.00624))
(32cc6532,CompactBuffer(0.02337, -0.05223, -0.03591))
(etc.)

一旦将它们分组,我就会尝试取平均值和标准差。我想简单地使用 .mean() 和 .sampleStdev()。当我尝试创建一个新的RDD的手段时,

val mean = grouped.mean()

返回错误

Error:(51, 22) value mean is not a member of org.apache.spark.rdd.RDD[(String, Iterable[Double])]

val mean = grouped.mean( )

我已经导入了 org.apache.spark.SparkContext._
我还尝试使用 sampleStdev()、.sum()、.stats() 得到相同的结果。不管是什么问题,它似乎都会影响所有数字 RDD 操作。

【问题讨论】:

标签: scala apache-spark


【解决方案1】:

让我们考虑以下几点:

val data = List(("32540b03",-0.00699), ("a93dec11",0.00624),
                ("32cc6532",0.02337) , ("32540b03",0.256023),
                ("32cc6532",-0.03591),("32cc6532",-0.03591))

val rdd = sc.parallelize(data.toSeq).groupByKey().sortByKey()

计算每对平均值的一种方法如下:

你需要定义一个平均方法:

def average[T]( ts: Iterable[T] )( implicit num: Numeric[T] ) = {
   num.toDouble( ts.sum ) / ts.size
}

您可以在 rdd 上应用您的方法,如下所示:

val avgs = rdd.map(x => (x._1, average(x._2)))

您可以检查:

avgs.take(3)

这就是结果:

res4: Array[(String, Double)] = Array((32540b03,0.1245165), (32cc6532,-0.016149999999999998), (a93dec11,0.00624))

【讨论】:

    【解决方案2】:

    官方的方法是使用reduceByKey 而不是groupByKey

    val result = sc.parallelize(data)
      .map { case (key, value) => (key, (value, 1)) }
      .reduceByKey { case ((value1, count1), (value2, count2))
        => (value1 + value2, count1 + count2)}
      .mapValues {case (value, count) =>  value.toDouble / count.toDouble}
    

    另一方面,您的解决方案中的问题是grouped(String, Iterable[Double]) 形式的对象的RDD(就像在错误中一样)。例如,您可以计算 Ints 或 doubles 的 RDD 的平均值,但成对的 rdd 的平均值是多少。

    【讨论】:

    • @ablcerek,如果您的键中有多个值,这会起作用吗?
    • @EB 我不确定我是否理解,但如果你想为元素的 rdds 计算它(key, value: List[Double]),那么你必须将第一个映射更改为{ case (key, value) => (key, (value.sum, value.count))}
    【解决方案3】:

    这是一个没有自定义函数的完整程序:

    val conf = new SparkConf().setAppName("means").setMaster("local[*]")
    val sc = new SparkContext(conf)
    
    val data = List(("Lily", 23), ("Lily", 50),
                    ("Tom", 66), ("Tom", 21), ("Tom", 69),
                    ("Max", 11), ("Max", 24))
    
    val RDD = sc.parallelize(data)
    
    val counts = RDD.map(item => (item._1, (1, item._2.toDouble)) )
    val countSums = counts.reduceByKey((x, y) => (x._1 + y._1, x._2 + y._2) )
    val keyMeans = countSums.mapValues(avgCount => avgCount._2 / avgCount._1)
    
    for ((key, mean) <- keyMeans.collect()) println(key + " " + mean)
    

    【讨论】:

      猜你喜欢
      • 2015-07-07
      • 1970-01-01
      • 2017-02-20
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多