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

PySpark多条件高级行过滤:贷款数据筛选需求及代码优化

问题

需要创建满足以下条件的新DataFrame:

  • 借款人(ID)至少有2笔贷款;
  • 每笔后续贷款金额小于前一笔,且在前一笔贷款结清后15天内发放;
  • 最后一笔未结清贷款(无ClosingDate)金额大于前一笔贷款。

现有DataFrame

Name    ID     ContractDate LoanSum ClosingDate
A       ID1    2022-10-10   10      2022-10-15 
A       ID1    2022-10-16   8       2022-10-25
A       ID1    2022-10-27   25      
B       ID2    2022-12-12   10      2022-10-15
B       ID2    2022-12-16   22      2022-11-18
B       ID2    2022-12-20   9       2022-11-25
B       ID2    2023-11-29   13      
C       ID3    2022-11-11   30      2022-11-18

预期结果

Name    ID     ContractDate LoanSum ClosingDate
A       ID1    2022-10-10   10      2022-10-15 
A       ID1    2022-10-16   8       2022-10-25
A       ID1    2022-10-27   25      
B       ID2    2022-12-12   10      2022-10-15
B       ID2    2022-12-16   22      2022-11-18
B       ID2    2022-12-20   9       2022-11-25
B       ID2    2023-11-29   13      

现有代码

用户尝试了以下代码,但无法正确筛选符合要求的未结清贷款:

cols = df.columns
w = Window.partitionBy('ID').orderBy('ContractDate')

newdf = df.withColumn('PreviousContractDate', f.lag('ContractDate').over(w)) \
  .withColumn('PreviousLoanSum', f.lag('LoanSum').over(w)) \
  .withColumn('Target', f.expr('datediff(ContractDate, PreviousContractDate) >= 1 and datediff(ContractDate, PreviousContractDate) < 16 and LoanSum - PreviousLoanSum < 0')) \
  .withColumn('Target', f.col('Target') | f.lead('Target').over(w)) \
  .filter('Target == True')

解决方案

要满足所有三个条件,需分步骤验证ID级的整体合规性,再保留符合要求的ID的所有记录,完整代码及解释如下:

完整代码

import pyspark.sql.functions as f
from pyspark.sql.window import Window

# 定义两个窗口:按ID分组排序、按ID分组聚合
w_sorted = Window.partitionBy('ID').orderBy('ContractDate')
w_id = Window.partitionBy('ID')

# 1. 计算每笔贷款的前置关联字段,标记各条件的合规性
df_processed = df.withColumn('prev_closing_date', f.lag('ClosingDate').over(w_sorted)) \
                 .withColumn('prev_loan_sum', f.lag('LoanSum').over(w_sorted)) \
                 # 标记当前贷款是否符合"后续贷<前一笔+结清后15天内发放"的条件
                 .withColumn('valid_sequence',
                             f.when(
                                 f.col('prev_closing_date').isNotNull(),
                                 (f.datediff(f.col('ContractDate'), f.col('prev_closing_date')) <= 15) &
                                 (f.col('LoanSum') < f.col('prev_loan_sum'))
                             ).otherwise(f.lit(False))) \
                 # 标记最后一笔未结清贷款是否符合"金额>前一笔"的条件
                 .withColumn('valid_last_unpaid',
                             f.when(
                                 (f.col('ClosingDate').isNull()) & 
                                 (f.row_number().over(w_sorted) == f.count('*').over(w_id)),
                                 f.col('LoanSum') > f.col('prev_loan_sum')
                             ).otherwise(f.lit(False))) \
                 # 计算每个ID的整体合规指标
                 .withColumn('has_valid_sequence', f.max('valid_sequence').over(w_id)) \
                 .withColumn('has_valid_last_unpaid',
                             f.when(
                                 f.count(f.when(f.col('ClosingDate').isNull(), 1)).over(w_id) > 0,
                                 f.max('valid_last_unpaid').over(w_id)
                             ).otherwise(f.lit(True))) \
                 .withColumn('loan_count', f.count('*').over(w_id))

# 2. 筛选符合所有条件的ID,保留原始列
final_df = df_processed.filter(
    (f.col('loan_count') >= 2) &
    (f.col('has_valid_sequence')) &
    (f.col('has_valid_last_unpaid'))
).select(df.columns)

final_df.show()

关键步骤说明

  1. 窗口区分:w_sorted用于按ID分组后按合同日期排序,处理单条记录的前置关联;w_id用于计算每个ID的整体统计值(如贷款总数、合规标记的最大值)。
  2. 条件2验证:通过lag获取前一笔贷款的结清日期,判断当前贷款是否在结清后15天内发放且金额更小,标记为valid_sequence。
  3. 条件3验证:定位每个ID的最后一笔未结清贷款,判断其金额是否大于前一笔;若ID无未结清贷款,该条件自动视为满足。
  4. ID级筛选:聚合每个ID的贷款总数、是否存在合规序列、是否满足未结清贷款条件,最终筛选出同时符合三个要求的ID,保留其所有原始记录。

内容的提问来源于stack exchange,提问作者lenpyspanacb

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 05:30:08