【问题标题】:pySpark find Median in a distributed way?pySpark以分布式方式查找中位数?
【发布时间】:2015-07-07 10:12:56
【问题描述】:

是否有可能以分布式方式找到火花中的中位数?我目前正在查找:SumAverageVarianceCount,使用以下代码:

dataSumsRdd = numRDD.filter(lambda x: filterNum(x[1])).map(lambda line: (line[0], float(line[1])))\
    .aggregateByKey((0.0, 0.0, 0.0),
     lambda (sum, sum2, count), value: (sum + value, sum2 + value**2, count+1.0),
     lambda (suma, sum2a, counta), (sumb, sum2b, countb): (suma + sumb, sum2a + sum2b, counta + countb))
#Generate RDD of Count, Sum, Average, Variance
dataStatsRdd = dataSumsRdd.mapValues(lambda (sum, sum2, count) : (count, sum, sum/count, round(sum2/count - (sum/count)**2, 7)))

我不太确定如何找到中位数。为了找到标准偏差,我只是用平方根方差在本地计算结果。一旦我收集到中值,我也可以轻松地在本地进行 Skewness。

我的数据在键/值对中(键 = 列)

【问题讨论】:

  • 看看this question。高效的分布式中值算法并不简单。

标签: apache-spark pyspark


【解决方案1】:

我正在看的是(它不是最好的方法......但我能想到的唯一方法):

def medianFunction(x):
    count = len(x)
    if count % 2 == 0:
        l = count / 2 - 1
        r = l + 1
        value = (x[l - 1] + x[r - 1]) / 2
        return value
    else:
        l = count / 2
        value = x[l - 1]
        return value

   medianRDD = numFilterRDD.groupByKey().map(lambda (x, y): (x, list(y))).mapValues(lambda x: medianFunction(x)).collect()

【讨论】:

  • 行 medianRDD = 以 .collect() 结尾。是故意的吗?您是否在一些测试数据上测试过这个解决方案?
  • .collect 是一个动作,它会产生对驱动程序没有危险的输出。你有什么顾虑?
猜你喜欢
  • 1970-01-01
  • 2014-09-20
  • 1970-01-01
  • 2019-11-14
  • 2020-09-12
  • 2011-12-04
  • 1970-01-01
  • 1970-01-01
  • 2021-01-17
相关资源
最近更新 更多