如何用PySpark窗口函数按指定type循环序列生成rank列?
如何用PySpark窗口函数生成指定循环序列的rank列?
原始DataFrame
A B type X 11 typeA X 12 typeA X 13 typeB X 14 typeB X 15 typeC X 16 typeC Y 17 typeA Y 18 typeA Y 19 typeB Y 20 typeB Y 21 typeC Y 22 typeC
需求说明
需要基于A列分区,生成rank列时遵循**typeA → typeB → typeC的循环序列**:同一分区内,先取第一个typeA,再第一个typeB,再第一个typeC;接着取第二个typeA,第二个typeB,第二个typeC,以此类推。期望结果如下:
期望结果
A B type rank X 11 typeA 1 X 13 typeB 2 X 15 typeC 3 X 12 typeA 4 X 14 typeB 5 X 16 typeC 6 Y 17 typeA 1 Y 19 typeB 2 Y 21 typeC 3 Y 18 typeA 4 Y 20 typeB 5 Y 22 typeC 6
当前错误实现
from pyspark.sql import Window result = (df .withColumn("rank", row_number() .over(Window.partitionBy('A') .orderBy(col('type').desc()) ) ) ).persist()
当前错误结果
A B type rank X 11 typeA 1 X 12 typeA 2 X 13 typeB 3 X 14 typeB 4 X 15 typeC 5 X 16 typeC 6 Y 17 typeA 1 Y 18 typeA 2 Y 19 typeB 3 Y 20 typeB 4 Y 21 typeC 5 Y 22 typeC 6
解决方案
要实现指定的循环序列rank,核心是先给每个分区内的同类型数据标记组内序号,再按组内序号+指定type顺序排序,最后生成rank:
from pyspark.sql import Window from pyspark.sql.functions import row_number, col, when # 1. 给每个(A, type)分组内的行分配组内序号(按B排序保证同type内的顺序与示例一致) window_group = Window.partitionBy("A", "type").orderBy("B") df_with_seq = df.withColumn("group_seq", row_number().over(window_group)) # 2. 定义type的排序优先级,确保typeA→typeB→typeC的顺序 df_with_priority = df_with_seq.withColumn( "type_priority", when(col("type") == "typeA", 1) .when(col("type") == "typeB", 2) .when(col("type") == "typeC", 3) ) # 3. 按A分区,先按组内序号排序,再按type优先级排序,生成最终rank window_rank = Window.partitionBy("A").orderBy("group_seq", "type_priority") result = df_with_priority.withColumn("rank", row_number().over(window_rank)).drop("group_seq", "type_priority") result.show()
逻辑说明
group_seq:标记同类型数据在分区内的批次(比如第一个typeA属于批次1,第二个属于批次2)type_priority:固定type的排序顺序,避免默认字符串排序的不确定性- 最终窗口按
批次序号+type优先级排序,生成的row_number就是符合要求的循环rank
内容的提问来源于stack exchange,提问作者Strayhorn
相关产品推荐
相关产品推荐

