PySpark复杂过滤操作:特定贷款场景数据筛选需求
PySpark实现贷款记录的复杂过滤
需求分析
需要筛选出符合以下条件的借款人(按唯一ID分组)的所有贷款记录:
- 借款人至少有两笔贷款
- 存在某一笔贷款未结清(
ClosingDate为null),且该笔贷款与下一笔已发放贷款的金额差≤5
代码实现
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql import functions as F # 初始化Spark会话 spark = SparkSession.builder.appName("LoanRecordFilter").getOrCreate() # 加载数据集(示例数据替换为你的实际数据源) data = [ ("A", "ID1", "2022-10-10", 10, "2022-10-15"), ("A", "ID1", "2022-10-15", 15, None), ("A", "ID1", "2022-10-30", 20, "2022-11-10"), ("B", "ID2", "2022-11-11", 15, "2022-10-14"), ("B", "ID2", "2022-12-10", 30, None), ("B", "ID2", "2022-12-12", 35, "2022-12-14"), ("C", "ID3", "2022-12-19", 19, "2022-11-10"), ("D", "ID4", "2022-12-10", 10, None), ("D", "ID4", "2022-12-12", 40, "2022-11-29") ] df = spark.createDataFrame(data, ["Name", "ID", "ContractDate", "LoanSum", "ClosingDate"]) # 转换日期字段类型,避免字符串排序异常 df = df.withColumn("ContractDate", F.to_date("ContractDate")) df = df.withColumn("ClosingDate", F.to_date("ClosingDate")) # 定义窗口:按ID分组,按合同日期升序排序 user_window = Window.partitionBy("ID").orderBy("ContractDate") user_count_window = Window.partitionBy("ID") # 添加辅助字段:行号、下一笔贷款金额、当前笔是否未结清、总贷款数 df_aux = df.withColumn("row_num", F.row_number().over(user_window)) \ .withColumn("next_loan_sum", F.lead("LoanSum").over(user_window)) \ .withColumn("is_unclosed", F.col("ClosingDate").isNull()) \ .withColumn("total_loans", F.count("*").over(user_count_window)) # 计算当前笔与下一笔的金额差绝对值 df_aux = df_aux.withColumn("amount_diff", F.abs(F.col("LoanSum") - F.col("next_loan_sum"))) # 筛选符合条件的借款人ID qualified_ids = df_aux.filter( (F.col("total_loans") >= 2) & (F.col("is_unclosed")) & (F.col("amount_diff") <= 5) ).select("ID").distinct() # 关联原始数据,保留符合条件的所有记录,并调整列名 result_df = df.join(qualified_ids, on="ID", how="inner") \ .withColumnRenamed("ClosingDate", "Status") # 查看结果 result_df.show()
代码说明
- 日期类型转换:将
ContractDate和ClosingDate转换为日期类型,确保排序逻辑正确。 - 窗口函数运用:
- 按ID分组并按合同日期排序,生成行号标记每笔贷款的顺序
- 使用
lead函数获取下一笔贷款的金额,用于计算金额差 - 标记当前贷款是否未结清,统计每个用户的总贷款数
- 条件筛选:筛选出满足“至少两笔贷款、存在未结清贷款且与下一笔金额差≤5”的用户ID
- 结果生成:通过关联筛选出的ID,保留这些用户的所有贷款记录,并将
ClosingDate重命名为Status以匹配预期输出。
内容的提问来源于stack exchange,提问作者lenpyspanacb
相关产品推荐
相关产品推荐

