You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.26 21:47:33