PySpark按部门分组后将分数除以组内最大值的实现方法
问题:按部门归一化Score值(除以组内最大值)
原始DataFrame结构如下:
columns = ['id', 'department', 'score'] vals = [ (1, 'AB', 141), (2, 'AB', 140), (3, 'AB', 210), (4, 'AB', 120), (5, 'EF', 20), (6, 'EF', 15) ]
需求:按department分组,计算每组score的最大值,再将组内所有score除以该最大值,最终得到类似如下结果(保留两位小数):
(1, 'AB', 0.67), (2, 'AB', 0.67), (3, 'AB', 1.00), (4, 'AB', 0.57), (5, 'EF', 1.00), (6, 'EF', 0.75)
目前尝试的代码:
>>> max_distance = df.groupby("department").agg({"score": "max"}).collect() >>> max_distance [Row(department='AB', max(score)=210.0), Row(department='EF', max(score)=20.0)]
解决方案
方法1:使用窗口函数(推荐,分布式高效)
窗口函数可以直接在原DataFrame上计算组内最大值,无需将数据拉到本地,适合大数据场景:
from pyspark.sql import Window import pyspark.sql.functions as F # 定义窗口规则:按department分区 dept_window = Window.partitionBy("department") # 计算组内最大值,再计算归一化后的分数并保留两位小数 result_df = df.withColumn("group_max_score", F.max("score").over(dept_window)) \ .withColumn("normalized_score", F.round(F.col("score") / F.col("group_max_score"), 2)) \ .select("id", "department", "normalized_score") # 输出结果 result_df.show()
执行后会得到符合预期的归一化结果,且全程在分布式集群中处理,避免了collect()带来的内存风险。
方法2:通过Join实现(基于已有分组结果)
如果要基于你已经得到的分组最大值结果,可以将其转为DataFrame后与原表关联,再计算归一化值:
import pyspark.sql.functions as F # 将分组结果转为DataFrame(避免collect,直接用groupBy后的结果) max_df = df.groupBy("department").agg(F.max("score").alias("group_max_score")) # 关联原DataFrame,计算归一化分数 result_df = df.join(max_df, on="department", how="inner") \ .withColumn("normalized_score", F.round(F.col("score") / F.col("group_max_score"), 2)) \ .select("id", "department", "normalized_score") result_df.show()
注意:尽量避免使用
collect()将分布式数据拉到Driver端,数据量较大时容易引发内存溢出,且处理效率远低于分布式操作。
内容的提问来源于stack exchange,提问作者theodre7
相关产品推荐
相关产品推荐

