如何使用Spark实现分组计数并输出(id, count)格式的元组结果
PySpark按ID分组统计实现方案
reduceByKey不是唯一实现方式,优先推荐使用DataFrame原生聚合API,性能比RDD的reduceByKey更高,Spark SQL层的Catalyst优化器会对执行计划做自动优化,数据量越大优势越明显。
方案1:DataFrame groupBy 聚合(最优方案)
直接基于加载好的DataFrame操作,无需转RDD,代码简洁高效:
from pyspark.sql.functions import count # 按ID分组统计出现次数,重命名结果列 result_df = myDF.groupBy("ID").agg(count("*").alias("COUNT")).orderBy("ID") # 预览结果 result_df.show() # 导出为CSV文件 # 小数据量可以加coalesce(1)合并为单文件输出,大数据量不建议加,避免单节点内存溢出 result_df.coalesce(1).write.format("csv").option("header", "true").save("filepath/to/save/result") # 获取(id, count)格式的元组列表 id_count_tuples = result_df.rdd.map(lambda row: (row.ID, row.COUNT)).collect()
方案2:RDD reduceByKey 实现(原有思路兼容方案)
如果必须使用RDD接口实现,可按如下方式编写,性能弱于方案1:
# 先将DataFrame转为( id, 1 )格式的键值对RDD,再按key聚合求和 result_rdd = myDF.rdd.map(lambda row: (row["ID"], 1)) \ .reduceByKey(lambda x, y: x + y) \ .sortByKey() # 转为DataFrame导出CSV result_df = result_rdd.toDF(["ID", "COUNT"]) result_df.write.format("csv").option("header", "true").save("filepath/to/save/result") # 直接获取元组列表 id_count_tuples = result_rdd.collect()
其他可选实现(不推荐)
还可以用groupByKey、aggregateByKey实现相同逻辑,但groupByKey会产生大量无效shuffle传输,性能远低于上述两种方案,不建议使用。
内容的提问来源于stack exchange,提问作者tomiyos654
相关产品推荐
相关产品推荐

