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

PySpark:如何在分区的子集内正确计算排名?

问题描述

尝试在PySpark中为分区子集(category='Y'的行)计算排名,原代码使用when语句筛选后调用rank(),但结果不符合预期:category='Y'且code=15的行rank为2,预期应为1。

原代码及执行结果如下:

import pyspark.sql.functions as F
from pyspark.sql.window import Window
from pyspark.sql import Row

data = [
    Row(id=1, code=14, category='N'),
    Row(id=1, code=20, category='Y'),
    Row(id=1,  code=19, category='Y'),
    Row(id=1,  code=22, category='Y'),
    Row(id=1,  code=15, category='Y'),
]

ps_df = spark.createDataFrame(data)

window = Window.partitionBy('id').orderBy('code')

ps_df = ps_df.withColumn('rank', F.when(F.col('category')=='Y', F.rank().over(window)))

ps_df.show()

执行结果:

+---+----+--------+----+
| id|code|category|rank|
+---+----+--------+----+
|  1|  14|       N|NULL|
|  1|  15|       Y|   2|
|  1|  19|       Y|   3|
|  1|  20|       Y|   4|
|  1|  22|       Y|   5|
+---+----+--------+----+
问题原因

原代码的窗口是基于整个id分区(包含category='N'的行)计算排名,when语句只是隐藏了非Y行的rank值,并没有改变rank的计算逻辑。code=14的N行在排序后占据了第1位,所以后续Y行的rank从2开始。

解决方案

有两种常用方法实现分区子集内的排名:

方法1:先筛选Y行计算排名,再关联回原表

先对category='Y'的行单独计算排名,再通过id和code关联到原表,非Y行的rank自动为null:

# 筛选Y行并计算排名
y_ranked = ps_df.filter(F.col('category') == 'Y') \
    .withColumn('rank', F.rank().over(Window.partitionBy('id').orderBy('code')))

# 关联回原表
result_df = ps_df.join(y_ranked, on=['id', 'code', 'category'], how='left')

result_df.show()

执行结果:

+---+----+--------+----+
| id|code|category|rank|
+---+----+--------+----+
|  1|  14|       N|NULL|
|  1|  15|       Y|   1|
|  1|  19|       Y|   2|
|  1|  20|       Y|   3|
|  1|  22|       Y|   4|
+---+----+--------+----+

方法2:使用条件窗口函数(通过累计计数实现排名)

在窗口中仅对category='Y'的行计数,计算当前行及之前的有效Y行数量,作为子集内的排名:

window = Window.partitionBy('id').orderBy('code')

# 计算当前行之前(含当前)的Y行数量,作为排名
result_df = ps_df.withColumn(
    'rank',
    F.when(
        F.col('category') == 'Y',
        F.sum(F.when(F.col('category') == 'Y', 1).otherwise(0)).over(window)
    )
)

result_df.show()

执行结果与方法1一致,这种方法无需额外关联操作,适合数据量较大的场景。

补充说明

如果需要处理并列排名(相同code的Y行排名相同),方法1中的rank()可直接满足;方法2则需调整为dense_rank逻辑,例如:

window = Window.partitionBy('id').orderBy(F.when(F.col('category')=='Y', F.col('code')))
result_df = ps_df.withColumn(
    'rank',
    F.when(F.col('category') == 'Y', F.dense_rank().over(window))
)

内容的提问来源于stack exchange,提问作者Henri

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 12:06:06