PySpark按Pick/Drop分组计算聚合值的正确实现方法
修正PySpark分组计算并合并结果的方案
原代码问题分析
- 同时按
Pick和Drop分组,导致FarePick、TipPick仅计算了当前Pick+Drop组合的总和,而非整个Pick分组的总费用/小费 AvgDragDrop与AvgDragPick计算逻辑完全一致,未实现按Drop分组的平均Drag值
修正方案
通过分别计算Pick、Drop分组的指标,再通过关联合并结果,满足需求:
from pyspark.sql import functions as func data = [ (1, 1, 4.00, 4.00, 1.00), (1, 2, 5.00, 10.00, 8.00), (1, 2, 5.00, 15.00, 12.00), (3, 2, 11.00, 12.00, 17.00), (3, 5, 41.00, 25.00, 13.00), (4, 6, 50.00, 70.00, 2.00) ] columns = ["Pick", "Drop", "Fare", "Tip", "Drag"] df = spark.createDataFrame(data, columns) # 1. 计算按Pick分组的总Fare、总Tip、平均Drag pick_agg = df.groupBy("Pick").agg( func.sum("Fare").alias("FarePick"), func.sum("Tip").alias("TipPick"), func.avg("Drag").alias("AvgDragPick") ) # 2. 计算按Drop分组的平均Drag drop_agg = df.groupBy("Drop").agg( func.avg("Drag").alias("AvgDragDrop") ) # 3. 合并结果:保留所有唯一的Pick-Drop组合,关联两个汇总表 result = df.select("Pick", "Drop").distinct() \ .join(pick_agg, on="Pick", how="inner") \ .join(drop_agg, on="Drop", how="inner") \ .orderBy("Pick") result.show()
可选:保留原数据所有行
如果需要保留原数据中重复的Pick-Drop记录,替换合并逻辑为:
result = df.join(pick_agg, on="Pick", how="inner") \ .join(drop_agg, on="Drop", how="inner") \ .select("Pick", "Drop", "FarePick", "TipPick", "AvgDragPick", "AvgDragDrop") \ .orderBy("Pick") result.show()
输出结果(去重后)
+----+----+--------+-------+-------------+-------------+ |Pick|Drop|FarePick|TipPick|AvgDragPick |AvgDragDrop | +----+----+--------+-------+-------------+-------------+ |1 |1 |14.0 |29.0 |7.0 |1.0 | |1 |2 |14.0 |29.0 |7.0 |12.3333333333| |3 |2 |52.0 |37.0 |15.0 |12.3333333333| |3 |5 |52.0 |37.0 |15.0 |13.0 | |4 |6 |50.0 |70.0 |2.0 |2.0 | +----+----+--------+-------+-------------+-------------+
内容的提问来源于stack exchange,提问作者Saffik
相关产品推荐
相关产品推荐

