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

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()

代码说明

  1. 日期类型转换:将ContractDate和ClosingDate转换为日期类型,确保排序逻辑正确。
  2. 窗口函数运用:
    • 按ID分组并按合同日期排序,生成行号标记每笔贷款的顺序
    • 使用lead函数获取下一笔贷款的金额,用于计算金额差
    • 标记当前贷款是否未结清,统计每个用户的总贷款数
  3. 条件筛选:筛选出满足“至少两笔贷款、存在未结清贷款且与下一笔金额差≤5”的用户ID
  4. 结果生成:通过关联筛选出的ID,保留这些用户的所有贷款记录,并将ClosingDate重命名为Status以匹配预期输出。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 13:55:32