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()
关键步骤说明
- 窗口区分:
w_sorted用于按ID分组后按合同日期排序,处理单条记录的前置关联;w_id用于计算每个ID的整体统计值(如贷款总数、合规标记的最大值)。 - 条件2验证:通过
lag获取前一笔贷款的结清日期,判断当前贷款是否在结清后15天内发放且金额更小,标记为valid_sequence。 - 条件3验证:定位每个ID的最后一笔未结清贷款,判断其金额是否大于前一笔;若ID无未结清贷款,该条件自动视为满足。
- ID级筛选:聚合每个ID的贷款总数、是否存在合规序列、是否满足未结清贷款条件,最终筛选出同时符合三个要求的ID,保留其所有原始记录。
内容的提问来源于stack exchange,提问作者lenpyspanacb
相关产品推荐
相关产品推荐

