PySpark多条件过滤:筛选符合特定规则的贷款DataFrame数据
问题描述
原始数据
Name ID ContractDate LoanSum ClosingDate A ID1 2022-10-10 10 2022-10-16 A ID1 2022-10-10 15 2022-10-18 A ID1 2022-10-20 20 2022-10-31 A ID1 2022-10-20 20 2022-10-30 A ID1 2022-11-10 14 2022-11-22 A ID1 2022-11-10 15 2022-11-22 B ID2 2022-11-11 15 2022-11-15 B ID2 2022-11-11 30 2022-11-18 B ID2 2022-11-17 35 2022-11-22 B ID2 2022-11-17 35 2022-11-24 C ID3 2022-12-19 19 2022-11-10
需求说明
按唯一ID分组后,筛选出符合以下条件的所有贷款记录:
- 存在至少两笔同日发放的贷款;
- 后续另有至少两笔同日发放的贷款,且后者的
ContractDate与前者任意一笔的ClosingDate间隔不超过5天。
预期输出
Name ID ContractDate LoanSum ClosingDate A ID1 2022-10-10 10 2022-10-16 A ID1 2022-10-10 15 2022-10-18 A ID1 2022-10-20 20 2022-10-31 A ID1 2022-10-20 20 2022-10-30
已尝试代码
from pyspark.sql import functions as f from pyspark.sql import Window df = spark.createDataFrame(data).toDF('Name','ID','ContractDate','LoanSum','ClosingDate') df.show() cols = df.columns w = Window.partitionBy('ID').orderBy('ContractDate') df.withColumn('PreviousContractDate', f.lag('ContractDate').over(w)) \ .withColumn('Target', f.expr('datediff(ContractDate, PreviousContractDate) == 0')) \ .withColumn('Target', f.col('Target') | f.lead('Target').over(w)) \ .filter('Target == True')
解决方案
实现思路
- 先筛选出所有属于「多笔同日发放」的记录,排除单条发放的批次;
- 对每个ID下的多笔发放日期批次,生成前后配对,检查后续批次的放款日与前批次任意还款日的间隔是否≤5天;
- 提取所有符合条件的日期批次对应的原始记录,得到最终结果。
完整代码
from pyspark.sql import functions as f from pyspark.sql import Window # 1. 筛选出所有多笔同日发放的候选记录 w_batch = Window.partitionBy("ID", "ContractDate") df_candidate = df.withColumn("batch_count", f.count("*").over(w_batch)) \ .filter(f.col("batch_count") >= 2) \ .drop("batch_count") # 2. 按ID和放款日期分组,收集该批次的所有还款日,并生成前后批次配对 batch_dates = df_candidate.groupBy("ID", "ContractDate") \ .agg(f.collect_list("ClosingDate").alias("closing_dates")) \ .orderBy("ID", "ContractDate") w_id = Window.partitionBy("ID").orderBy("ContractDate") batch_pairs = batch_dates.withColumn("next_contract_date", f.lead("ContractDate").over(w_id)) \ .filter(f.col("next_contract_date").isNotNull()) # 3. 检查后续批次与前批次的还款日间隔是否≤5天 batch_pairs_expanded = batch_pairs.withColumn("closing_date", f.explode("closing_dates")) \ .withColumn("days_diff", f.datediff(f.col("next_contract_date"), f.col("closing_date"))) \ .filter((f.col("days_diff") >= 0) & (f.col("days_diff") <= 5)) \ .select("ID", "ContractDate", "next_contract_date") \ .distinct() # 4. 合并符合条件的前后日期批次,关联原始数据得到最终结果 valid_dates = batch_pairs_expanded.select("ID", "ContractDate") \ .union(batch_pairs_expanded.select("ID", "next_contract_date").withColumnRenamed("next_contract_date", "ContractDate")) \ .distinct() final_df = df_candidate.join(valid_dates, on=["ID", "ContractDate"], how="inner") final_df.show()
代码说明
- 第一步通过窗口函数统计每个
(ID, ContractDate)组的记录数,筛选出记录数≥2的批次,得到候选数据集; - 第二步将候选数据按批次分组,收集每个批次的还款日,再通过
lead函数获取每个批次的下一个多笔发放批次; - 第三步展开每个批次的还款日,计算与下一批次放款日的间隔,筛选出间隔在0到5天内的有效配对;
- 最后合并有效配对中的前后日期,关联候选数据集,得到所有符合条件的贷款记录。
内容的提问来源于stack exchange,提问作者lenpyspanacb
相关产品推荐
相关产品推荐

