如何用单行PySpark RDD操作实现元素求和后计算平均值
PySpark RDD单行计算平均值方案
你可以直接合并现有逻辑得到最简单行写法,也可以采用性能更优的单次RDD扫描写法:
最简合并写法(兼容你原有的实现逻辑)
直接将分步操作链式调用即可,无需额外导包:
result = rdd_example.map(lambda x: x[1]).reduce(lambda a, b: a + b) / rdd_example.count()
如果要沿用你之前导入的add算子,写法如下:
from operator import add result = rdd_example.map(lambda x: x[1]).reduce(add) / rdd_example.count()
注意:该写法会触发2次RDD作业,小数据量下无感知,大数据量场景下推荐下方的高性能写法。
额外提示:你给出的示例RDD中姓名属于字符串类型,实际定义时需要加引号,正确写法为rdd_example = sc.parallelize([("eliana",1),("peter",2),("andrew",3),("paul",4),("jhon",5)])
高性能单次扫描写法(大数据量场景推荐)
使用aggregate算子在一次RDD遍历中同时完成求和和计数,仅触发1次作业,性能更高,还自带除0保护避免空RDD报错:
result = rdd_example.aggregate( (0, 0), lambda acc, val: (acc[0] + val[1], acc[1] + 1), lambda acc1, acc2: (acc1[0] + acc2[0], acc1[1] + acc2[1]) ).let(lambda res: res[0]/res[1] if res[1] > 0 else 0)
参数说明:
- 第一个参数
(0, 0)为初始累加元组,首位存储数值求和结果,次位存储元素计数 - 第二个参数定义分区内的计算规则:每遍历一个RDD元素,就更新求和值和计数值
- 第三个参数定义分区间的合并规则:将不同分区计算得到的累加结果合并
- 末尾
let算子对最终的累加元组做除法运算,同时加入了空RDD的除0保护逻辑
内容的提问来源于stack exchange,提问作者Andres
相关产品推荐
相关产品推荐

