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

PySpark如何按Matricule分组筛选每个分组内Value最高的所有行

PySpark按分组取最大值所有记录的实现方法

实现方案1:窗口函数(推荐)

这是性能最优、代码最简洁的实现方式:

  1. 构造按Matricule分区、Value降序排序的窗口
  2. 用rank()函数给窗口内记录排名,相同Value的记录排名一致
  3. 筛选排名为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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 09:15:03