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

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')
解决方案

实现思路

  1. 先筛选出所有属于「多笔同日发放」的记录,排除单条发放的批次;
  2. 对每个ID下的多笔发放日期批次,生成前后配对,检查后续批次的放款日与前批次任意还款日的间隔是否≤5天;
  3. 提取所有符合条件的日期批次对应的原始记录,得到最终结果。

完整代码

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 21:27:53