PySpark中RDD对象按字符串元素分组并聚合其余数值的实现方法
PySpark RDD分组求和实现方案
你可以直接用reduceByKey实现,性能比groupByKey更好,能避免中间shuffle阶段的冗余开销,完整实现代码如下:
# 假设你原有RDD的变量名为origin_rdd # 第一步:将RDD转换为(分组键, 数值元组)的键值对格式 key_value_rdd = origin_rdd.map(lambda x: (x[0], x[1:])) # 第二步:按键分组,对应位置的数值分别累加求和 sum_rdd = key_value_rdd.reduceByKey(lambda a, b: tuple(a_i + b_i for a_i, b_i in zip(a, b))) # 如需查看结果可执行动作算子 print(sum_rdd.take(3))
实现逻辑说明
map操作里把每个元组的第一个字符串元素设为分组键,提取剩余的所有数值作为该键对应的值reduceByKey操作中,对同一个分组键下的每两组数值元组,通过zip逐位配对后相加,最终得到该分组下所有数值的累加结果
如果确实需要用groupByKey实现(性能更低,仅做参考),可参考以下写法:
sum_rdd = origin_rdd.map(lambda x: (x[0], x[1:])) \ .groupByKey() \ .mapValues(lambda values: tuple(sum(col) for col in zip(*values)))
内容的提问来源于stack exchange,提问作者Rashida
相关产品推荐
相关产品推荐

