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

如何利用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))

方法一:窗口函数实现

你当前代码的核心问题:

  1. 重复使用同一个变量名w_max_rnk,覆盖了之前定义的窗口
  2. 直接引用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 23:25:27