PySpark如何按Matricule分组筛选每个分组内Value最高的所有行
PySpark按分组取最大值所有记录的实现方法
实现方案1:窗口函数(推荐)
这是性能最优、代码最简洁的实现方式:
- 构造按
Matricule分区、Value降序排序的窗口 - 用
rank()函数给窗口内记录排名,相同Value的记录排名一致 - 筛选排名为1的所有记录即可
完整代码如下:
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化Spark spark = SparkSession.builder.appName("test").getOrCreate() # 构造测试数据(对应你提供的原始表) test_data = [ ("MA12", "A101", 25), ("MA12", "K215", 25), ("MA12", "C231", 70), ("MA12", "G348", 70), ("MA12", "B401", 70), ("MA12", "E291", 70), ("MA20", "D34", 16), ("MA20", "A45", 16), ("MA20", "A40", 15), ("MA20", "G16", 18), ("MA20", "K26", 18) ] df = spark.createDataFrame(test_data, ["Matricule", "Count", "Value"]) # 核心逻辑 window_spec = Window.partitionBy("Matricule").orderBy(F.desc("Value")) result = df.withColumn("rank", F.rank().over(window_spec)) \ .filter(F.col("rank") == 1) \ .drop("rank") # 输出结果 result.show()
实现方案2:分组聚合后Join
适合超大规模数据集场景,逻辑如下:
# 先计算每个分组的最大Value max_df = df.groupBy("Matricule").agg(F.max("Value").alias("max_val")) # 关联原表过滤出符合条件的记录 result = df.join(max_df, (df.Matricule == max_df.Matricule) & (df.Value == max_df.max_val)) \ .select(df["*"])
两种方案的输出都完全匹配你给出的预期结果。
内容的提问来源于stack exchange,提问作者elokema
相关产品推荐
相关产品推荐

