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
相关产品推荐
相关产品推荐

