PySpark RDD reduceByKey()使用问题:无法按Key聚合数值与计数
PySpark代码问题排查与修正
问题根源
你的代码中map(lambda x: (x,1))把整个原元组(比如("a",10))作为了key,导致每个元素的key都是唯一的——("a",10)和("a",8)是两个不同的key,reduceByKey自然无法对相同字母的元素进行分组聚合,最终输出和原数据几乎一致,没有完成求和统计。
修正方案
需要调整map阶段的key设置,以原元组的第一个元素(字母)作为分组key,同时把数值和计数1打包成value,后续再通过reduceByKey聚合:
完整修正代码
rdd_avg=sc.parallelize([("a",10),("b",15),("c",20),("a",8)]) # 用字母作为分组key,(数值, 计数)作为value rdd_sp1=rdd_avg.map(lambda x: (x[0], (x[1], 1))) # 按key聚合:数值求和,计数求和 rdd_sp2=rdd_sp1.reduceByKey(lambda a,b:(a[0]+b[0], a[1]+b[1])) # 转换为期望的输出格式 rdd_final = rdd_sp2.map(lambda x: ((x[0], x[1][0]), x[1][1])) print(rdd_final.collect())
输出结果
[(('c', 20), 1), (('a', 18), 2), (('b', 15), 1)]
内容的提问来源于stack exchange,提问作者Dwaipayan Sarkar
相关产品推荐
相关产品推荐

