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

PySpark高级过滤:筛选符合条件的借贷记录

PySpark实现同ID同LoanSum且日期间隔小于5天的记录筛选

解决思路

核心是通过窗口函数对同一ID和LoanSum分组内的记录按日期排序,计算相邻记录的日期间隔,最终筛选出存在相邻日期间隔小于5天的所有相关记录。

完整实现代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, datediff, lag, lead
from pyspark.sql.window import Window
from pyspark.sql.types import DateType

# 初始化SparkSession(按需执行)
spark = SparkSession.builder.appName("LoanRecordFilter").getOrCreate()

# 1. 转换日期字段类型(原始数据中ContractDate为字符串,需转为日期类型)
df = df.withColumn("ContractDate", col("ContractDate").cast(DateType()))

# 2. 定义窗口规则:按ID、LoanSum分区,按ContractDate升序排序
window_spec = Window.partitionBy("ID", "LoanSum").orderBy("ContractDate")

# 3. 计算当前记录与前后记录的日期差
df_with_diffs = df.withColumn(
    "prev_date_diff", 
    datediff(col("ContractDate"), lag("ContractDate").over(window_spec))
).withColumn(
    "next_date_diff", 
    datediff(lead("ContractDate").over(window_spec), col("ContractDate"))
)

# 4. 筛选符合条件的记录:存在相邻日期间隔小于5天
result_df = df_with_diffs.filter(
    (col("prev_date_diff") < 5) | (col("next_date_diff") < 5)
).select("Name", "ID", "ContractDate", "LoanSum", "Status")

# 查看结果
result_df.show()

代码说明

  • 日期类型转换:确保ContractDate是日期类型,才能用datediff计算间隔。
  • 窗口函数:partitionBy("ID", "LoanSum")保证只在同一借款人同金额的分组内处理;orderBy("ContractDate")让记录按时间排序。
  • 相邻日期差计算:
    • lag获取前一条记录的日期,计算当前日期与前一条的间隔;
    • lead获取后一条记录的日期,计算后一条与当前日期的间隔。
  • 筛选逻辑:只要当前记录与前/后记录的间隔小于5天,就保留该记录,这样就能把所有符合“同一ID同LoanSum且存在间隔小于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.26 05:43:19