PySpark partitionBy窗口函数未正常工作,求错误排查
问题排查:Spark窗口函数行号计算不符合预期
原始数据表
| COL1 | COL2 | COL3 |
|---|---|---|
| COMP | 0005 | 2008-08-04 |
| COMP | 0009 | 2002-01-01 |
| COMP | 01.0 | 2002-01-01 |
| COMP | 0005 | 2008-01-01 |
| COMP | 0005 | 2001-10-20 |
| CTEC | 0009 | 2001-10-20 |
| COMP | 0005 | 2009-10-01 |
| COMP | 01.0 | 2003-07-01 |
| COMP | 02.0 | 2004-01-01 |
| CTEC | 0009 | 2021-09-24 |
需求
按COL1和COL2联合分区,对每个分区内的COL3降序排序后添加行号。
错误代码
windowSpec = Window.partitionBy(col("COL1")).partitionBy(col("COl2")).orderBy(desc("COL3")) TBL = TBL.withColumn(f"RANK", F.row_number().over(windowSpec))
问题现象
预期输出中COMP的COL2=0009行号应为1,CTEC的COL2=0009行号应为1,但实际输出中前者行号为2,后者为3。
错误原因及修复方案
1. 分区逻辑错误
连续调用两次partitionBy会覆盖之前的分区配置,而非叠加。你的代码最终仅以COL2作为分区列,导致COMP和CTEC中COL2=0009的行被分到同一个分区,行号自然连续递增。
2. 列名拼写错误
代码中col("COl2")的字母L是小写,与数据表中的COL2(大写L)不匹配,可能导致分区逻辑进一步异常。
正确代码
将多个分区列放在同一个partitionBy调用中,用逗号分隔:
windowSpec = Window.partitionBy(col("COL1"), col("COL2")).orderBy(desc("COL3")) TBL = TBL.withColumn("RANK", F.row_number().over(windowSpec))
验证说明
修正后的窗口会先按COL1分区,再在每个COL1分区内按COL2细分,每个子分区内按COL3降序生成行号,完全符合你的预期输出。
内容的提问来源于stack exchange,提问作者Debtanu Gupta
相关产品推荐
相关产品推荐

