基于PySpark实现分组求和金额字段、其余字段取组内最后值的方法咨询
PySpark混合聚合场景最优实现方案
无需拆分两张表再做关联,以下两种方案都仅需一次shuffle即可完成计算,性能远高于多轮分组+关联的实现。
前置说明
你需要先明确组内「最后一条记录」的排序规则:
- 若有业务时间类字段,优先用业务字段排序保证结果符合业务逻辑
- 若仅需按数据默认存储顺序排序,可先执行以下代码生成行序ID:
from pyspark.sql import functions as F df = df.withColumn("row_id", F.monotonically_increasing_id())
方案1:单轮groupBy聚合(性能最优)
全程仅触发一次分组shuffle,资源消耗最低,适合大数据量场景:
from pyspark.sql import functions as F result_df = df.orderBy("row_id") \ .groupBy("Composed_PKey") \ .agg( # 数值字段直接求和,sum原生忽略null值 F.sum("Auftragsmenge").alias("Auftragsmenge"), # 非统计字段取最后一个非空有效值 F.last("Vertriebsgrund_Key", ignorenulls=True).alias("Vertriebsgrund_Key"), F.last("Retoure_Soll_Haben_Kennzeichen", ignorenulls=True).alias("Retoure_Soll_Haben_Kennzeichen") )
方案2:窗口函数实现(逻辑直观,易调试)
适合需要中间结果排查问题的场景,性能略优于多轮分组,逻辑更易理解:
from pyspark.sql import Window from pyspark.sql import functions as F # 定义窗口:按主键分区,按行ID排序,帧覆盖整个分组 w = Window.partitionBy("Composed_PKey") \ .orderBy("row_id") \ .rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing) result_df = df.withColumn("Auftragsmenge", F.sum("Auftragsmenge").over(w)) \ .withColumn("Vertriebsgrund_Key", F.last("Vertriebsgrund_Key", ignorenulls=True).over(w)) \ .withColumn("Retoure_Soll_Haben_Kennzeichen", F.last("Retoure_Soll_Haben_Kennzeichen", ignorenulls=True).over(w)) \ .select("Composed_PKey", "Vertriebsgrund_Key", "Retoure_Soll_Haben_Kennzeichen", "Auftragsmenge") \ .dropDuplicates(["Composed_PKey"])
结果验证
用你提供的样例数据执行以上代码,最终输出结果如下:
| Composed_PKey | Vertriebsgrund_Key | Retoure_Soll_Haben_Kennzeichen | Auftragsmenge |
|---|---|---|---|
| 1016626879-000010-ST | 33 | X | 40.0 |
完全符合你的聚合规则要求。
内容的提问来源于stack exchange,提问作者Mow.Massimo
相关产品推荐
相关产品推荐

