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
相关产品推荐
相关产品推荐

