PySpark如何获取指定列最大值对应的包含所有列的整行记录
代码错误原因
你写的代码错误原因是将Name、Score、Section三个列共同作为分组键,Spark分组时会把这三个列值完全相同的行归为同一组,你的样例数据里每行的三个列值组合都是唯一的,相当于每行单独成组,聚合得到的max(Score)就是当前行自身的Score,自然无法筛选出全局最高分对应的行。
可行实现方案
方案1:先取最大值再过滤(最简单直观)
适合只需要取全局最大值对应行的场景,代码逻辑最简洁:
from pyspark.sql.functions import max # 计算全局最高Score值 global_max_score = df.select(max("Score")).first()[0] # 过滤出Score等于全局最大值的所有行 dataframe_max = df.filter(df.Score == global_max_score)
方案2:窗口函数实现(扩展性更强)
如果后续需要调整为按分组(比如按Section取每个分组的最高分),只需要修改窗口定义即可,不需要大改逻辑:
from pyspark.sql.functions import rank from pyspark.sql.window import Window # 定义窗口规则:全局按Score降序排序,如果要按Section分组取每组最高,添加.partitionBy("Section")即可 window_spec = Window.orderBy(df.Score.desc()) # 给每行计算排名,Score相同的行排名相同 df_with_rank = df.withColumn("score_rank", rank().over(window_spec)) # 筛选排名第一的行,删除辅助排名列 dataframe_max = df_with_rank.filter(df_with_rank.score_rank == 1).drop("score_rank")
方案3:基于groupby+join实现
如果你坚持要用groupby聚合的思路,可以先聚合得到最大值,再关联回原表得到完整行:
from pyspark.sql.functions import max # 聚合得到全局最大Score max_score_df = df.agg(max("Score").alias("max_score")) # 关联原表筛选出Score等于最大值的行 dataframe_max = df.join(max_score_df, df.Score == max_score_df.max_score, how="inner").drop("max_score")
内容的提问来源于stack exchange,提问作者Jesse M
相关产品推荐
相关产品推荐

