You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何获取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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.02 02:06:11