如何利用PySpark窗口函数为DataFrame添加最大、最小排名对应进程列?
问题描述
我有一个包含Process(进程名称)和Process_rank(进程排名)两列的PySpark DataFrame,希望新增两列,将最大排名、最小排名对应的进程名称展示在每一行中。参考示例输出列Max Rank Process (output I want using windowing)和Min Rank Process (output I want using windowing 2)查看预期结果。尝试用窗口函数时发现无法直接引用无聚合的列,求可行实现方案(无论是否用窗口函数)。
附当前尝试的代码:
from pyspark.sql.types import StructType,StructField, StringType, IntegerType from pyspark.sql import functions as F from pyspark.sql.window import Window schema = StructType([ \ StructField("Process",StringType(),True), \ StructField("Process_rank",IntegerType(),True), \ StructField("Max Rank Process (output I want using windowing)",StringType(),True) , \ StructField("Min Rank Process (output I want using windowing 2)",StringType(),True) ]) data = [("Inventory", 1, "Retire","Inventory"), \ ("Data availability", 2, "Retire", "Inventory"), \ ("Code Conversion", 3, "Retire", "Inventory"), \ ("Retire", 4, "Retire", "Inventory") ] df = spark.createDataFrame(data=data,schema=schema) ############Partitions # window1: partition by Process name, order by rank max w_max_rnk = Window.partitionBy("Process").orderBy(F.col("Process_rank").desc()) # window2: partition by Process name, order by rank min w_max_rnk = Window.partitionBy("Process").orderBy(F.col("Process_rank").asc()) #windowed cols to find max and min processes from dataframe df = df.withColumn("max_ranked_process", F.col("Process").over(w_max_rnk)) \ .withColumn("min_ranked_process", F.col("Process").over(w_max_rnk))
方法一:窗口函数实现
你当前代码的核心问题:
- 重复使用同一个变量名
w_max_rnk,覆盖了之前定义的窗口 - 直接引用
Process列加窗口无法自动定位到极值行,需要结合first()/last()聚合函数,同时窗口不需要按Process分区(若需求是全局的最大/最小排名进程)
修正后的代码:
# 定义全局窗口:不分区,按排名排序后覆盖全量数据行 w_global_max = Window.orderBy(F.col("Process_rank").desc()).rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing) w_global_min = Window.orderBy(F.col("Process_rank").asc()).rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing) # 用first()获取排序后的首行,即对应极值的进程名 df_result = df.withColumn( "max_ranked_process", F.first("Process").over(w_global_max) ).withColumn( "min_ranked_process", F.first("Process").over(w_global_min) ) df_result.show()
如果需要按某个字段分组计算每组的极值进程,只需在窗口中添加partitionBy("分组字段")即可。
方法二:非窗口函数(聚合+关联)实现
先计算全局的最大、最小排名进程名,再与原DataFrame关联,让每行都显示这两个值:
方式1:直接聚合取值
# 计算全局最大排名对应的进程 max_rank_proc = df.groupBy().agg( F.first("Process").orderBy(F.col("Process_rank").desc()).alias("max_ranked_process") ).collect()[0]["max_ranked_process"] # 计算全局最小排名对应的进程 min_rank_proc = df.groupBy().agg( F.first("Process").orderBy(F.col("Process_rank").asc()).alias("min_ranked_process") ).collect()[0]["min_ranked_process"] # 给原表新增列 df_result = df.withColumn("max_ranked_process", F.lit(max_rank_proc)) \ .withColumn("min_ranked_process", F.lit(min_rank_proc)) df_result.show()
方式2:关联聚合表(避免collect)
# 生成包含全局极值进程的聚合表 agg_df = df.groupBy().agg( F.first("Process").orderBy(F.col("Process_rank").desc()).alias("max_ranked_process"), F.first("Process").orderBy(F.col("Process_rank").asc()).alias("min_ranked_process") ) # 交叉关联原表,让每行都带上极值进程名 df_result = df.crossJoin(agg_df) df_result.show()
内容的提问来源于stack exchange,提问作者SPena
相关产品推荐
相关产品推荐

