如何在PySpark中筛选指定分组内的最高与最低rank值数据?
在PySpark中实现按分组筛选rank最值对应行的功能
需求:按id和hour分组,筛选出每组内rank值最小和最大对应的行,功能等效于以下Pandas代码:
参考Pandas代码
df.loc[ df.groupby(["id", "hour"])["rank"] \ .agg(["idxmin", "idxmax"]) \ .stack() ].sort_index()
输入数据
id year month date hour minute rank 54807 2021 12 31 6 29 1.0 54807 2021 12 31 6 31 2.0 54807 2021 12 31 7 15 1.0 54807 2021 12 31 7 18 2.0 54807 2021 12 31 7 30 3.0
期望输出
id year month date hour minute rank 54807 2021 12 31 6 29 1.0 54807 2021 12 31 6 31 2.0 54807 2021 12 31 7 15 1.0 54807 2021 12 31 7 30 3.0
PySpark实现方案
通过窗口函数计算分组内的rank极值,再筛选匹配行即可实现需求:
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import col, min, max # 初始化Spark会话 spark = SparkSession.builder.appName("rank_min_max_filter").getOrCreate() # 构建输入DataFrame data = [ (54807, 2021, 12, 31, 6, 29, 1.0), (54807, 2021, 12, 31, 6, 31, 2.0), (54807, 2021, 12, 31, 7, 15, 1.0), (54807, 2021, 12, 31, 7, 18, 2.0), (54807, 2021, 12, 31, 7, 30, 3.0) ] columns = ["id", "year", "month", "date", "hour", "minute", "rank"] df = spark.createDataFrame(data, columns) # 定义分组窗口:按id和hour分区 window_spec = Window.partitionBy("id", "hour") # 添加分组内的最小rank和最大rank字段 df_with_extremes = df.withColumn("min_rank", min("rank").over(window_spec)) \ .withColumn("max_rank", max("rank").over(window_spec)) # 筛选rank等于组内最小值或最大值的行 result = df_with_extremes.filter((col("rank") == col("min_rank")) | (col("rank") == col("max_rank"))) # 排序后输出结果 result.orderBy("id", "hour", "rank").show()
代码说明
- 用
Window.partitionBy("id", "hour")定义分组逻辑,对应Pandas的groupby(["id", "hour"]) - 通过窗口函数
min和max计算每组的rank极值 - 筛选rank值匹配极值的行,得到目标结果
- 按
id、hour、rank排序,对齐Pandas的sort_index()效果
内容的提问来源于stack exchange,提问作者Nabih Bawazir
相关产品推荐
相关产品推荐

