如何在PySpark DataFrame中按多条件为每行设置行号?
解决PySpark按条件生成分组行号的问题
原始DataFrame结构
id_company id_client id_loan date c1 id1 m1 2024-10-15 c2 id1 m2 2024-10-16 c3 id2 m3 2024-10-18 c3 id2 m4 2024-10-18 c3 id2 m5 2024-10-18 c4 id3 m6 2024-10-19 c4 id3 m7 2024-10-20 c4 id3 m8 2024-10-20 c5 id4 m9 2024-10-30 c6 id4 m10 2024-10-31 c6 id4 m11 2024-10-31
需求说明
- 以
id_client为分组维度,为每组内的记录生成行号 - 同一
id_client下,**不同id_company或不同date**的记录,行号依次递增 - 同一
id_client、同一id_company且同一date的记录,行号保持一致
预期结果
id_company id_client id_loan date row_number c1 id1 m1 2024-10-15 1 c2 id1 m2 2024-10-16 2 c3 id2 m3 2024-10-18 1 c3 id2 m4 2024-10-18 1 c3 id2 m5 2024-10-18 1 c4 id3 m6 2024-10-19 1 c4 id3 m7 2024-10-20 2 c4 id3 m8 2024-10-20 2 c5 id4 m9 2024-10-30 1 c6 id4 m10 2024-10-31 2 c6 id4 m11 2024-10-31 2
当前代码问题
原代码仅按id_subject(疑似笔误,应为id_client)分区、按date排序使用row_number(),但row_number()会给同分区同排序键的每条记录分配唯一递增号,且未纳入id_company作为分组判断维度,无法满足同id_client+id_company+date共享行号的需求:
from pyspark.sql import Window from pyspark.sql.functions import row_number df1 = df.withColumn("row_num", row_number().over(Window.partitionBy("id_subject").orderBy('date')))
解决方案
使用dense_rank()替代row_number(),并在窗口排序中同时纳入id_company和date,确保相同组合的记录共享行号:
from pyspark.sql import Window from pyspark.sql.functions import dense_rank # 定义窗口:按id_client分区,按id_company和date排序 window_spec = Window.partitionBy("id_client").orderBy("id_company", "date") # 生成符合要求的行号 df_result = df.withColumn("row_number", dense_rank().over(window_spec)) # 输出结果 df_result.show()
关键说明
dense_rank()的特性:同一分区内,排序键完全相同的记录会得到相同的排名,且排名连续无间隔,完美匹配需求中同组合共享行号、不同组合依次递增的规则。- 窗口的
orderBy必须包含id_company和date两个字段,才能准确区分不同的记录组合,确保行号分配符合规则。
内容的提问来源于stack exchange,提问作者lenpyspanacb
相关产品推荐
相关产品推荐

