如何获取PySpark DataFrame中每个ID首笔逾期贷款的发行索引
问题:获取首笔逾期贷款的发行索引序号
原始DataFrame如下:
Name ID ContractDate LoanSum DurationOfDelay A ID1 2023-01-01 10 10 A ID1 2023-01-03 15 15 A ID1 2022-12-29 20 0 A ID1 2022-12-28 40 0 B ID2 2023-01-05 15 19 B ID2 2023-01-10 30 0 B ID2 2023-01-07 35 25 B ID2 2023-01-06 35 0
目标是为每个唯一ID(或Name)展示首笔DurationOfDelay > 0的贷款的发行索引序号,期望结果:
Name ID IndexNum A ID1 3 B ID2 1
说明:
- 对于ID1,按发行日期排序后顺序为2022-12-28、2022-12-29、2023-01-01、2023-01-03,首笔逾期贷款是第3笔
- 对于ID2,按发行日期排序后顺序为2023-01-05、2023-01-06、2023-01-07、2023-01-10,首笔逾期贷款是第1笔
现有代码能筛选出首笔逾期贷款,但无法获取发行索引序号:
window_spec_subset = Window.partitionBy('ID').orderBy('ContractDate') subset = df.filter(F.col('DurationOfDelay') > 0) \ .withColumn('row_num', F.row_number().over(window_spec_subset)) \ .filter(F.col('row_num') == 1) \ .drop('row_num') subset.show()
运行结果:
+----+---+------------+-------+---------------+ |Name| ID|ContractDate|LoanSum|DurationOfDelay| +----+---+------------+-------+---------------+ | A|ID1| 2023-01-01| 10| 10| | B|ID2| 2023-01-05| 15| 19| +----+---+------------+-------+---------------+
解决方案
核心思路是先给每个用户的所有贷款按发行日期生成全局索引序号,再筛选首笔逾期贷款并提取对应索引,而非先筛选逾期记录再排序。
方法一:分步处理
from pyspark.sql import functions as F from pyspark.sql.window import Window # 1. 给每个用户的所有贷款按发行日期生成从1开始的全局索引 window_all = Window.partitionBy('ID').orderBy('ContractDate') df_with_index = df.withColumn('IndexNum', F.row_number().over(window_all)) # 2. 筛选逾期记录,再取每个用户的首笔逾期 window_overdue = Window.partitionBy('ID').orderBy('ContractDate') result = df_with_index.filter(F.col('DurationOfDelay') > 0) \ .withColumn('overdue_row_num', F.row_number().over(window_overdue)) \ .filter(F.col('overdue_row_num') == 1) \ .select('Name', 'ID', 'IndexNum') result.show()
运行结果:
+----+---+--------+ |Name| ID|IndexNum| +----+---+--------+ | A|ID1| 3| | B|ID2| 1| +----+---+--------+
方法二:用聚合函数简化代码
生成全局索引后,直接通过first聚合函数提取首笔逾期的索引:
from pyspark.sql import functions as F from pyspark.sql.window import Window window_all = Window.partitionBy('ID').orderBy('ContractDate') df_with_index = df.withColumn('IndexNum', F.row_number().over(window_all)) # 聚合获取每个ID首笔逾期的IndexNum result = df_with_index.filter(F.col('DurationOfDelay') > 0) \ .groupBy('Name', 'ID') \ .agg(F.first('IndexNum').alias('IndexNum')) result.show()
此方法同样能得到符合预期的结果,代码更简洁。
内容的提问来源于stack exchange,提问作者lenpyspanacb
相关产品推荐
相关产品推荐

