PySpark高级过滤需求:筛选符合条件的贷款记录
PySpark 高级过滤解决方案
针对你的需求,我们可以通过窗口函数计算相邻记录差值 + 关联匹配有效记录的方式实现,具体步骤如下:
1. 数据预处理:转换日期类型
首先将字符串格式的ContractDate转换为Spark日期类型,方便后续计算日期差:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 转换日期列 df = df.withColumn("ContractDate", F.to_date("ContractDate", "yyyy-MM-dd"))
2. 计算相邻记录的差值
按ID分组、ContractDate排序,使用lag函数获取上一条记录的日期和贷款额,计算两者的差值:
# 定义窗口:按ID分区,按合同日期升序排序 window_spec = Window.partitionBy("ID").orderBy("ContractDate") # 添加上一条记录的日期、贷款额,以及差值计算列 df_with_diff = df.withColumn("prev_date", F.lag("ContractDate").over(window_spec)) \ .withColumn("prev_loan_sum", F.lag("LoanSum").over(window_spec)) \ .withColumn("days_diff", F.datediff("ContractDate", "prev_date")) \ .withColumn("sum_diff", F.abs(F.col("LoanSum") - F.col("prev_loan_sum")))
3. 筛选满足条件的记录对
过滤出日期间隔<15天且贷款额差值≤3的记录对,然后提取这些记录的ID和ContractDate作为有效标识:
# 筛选符合条件的记录对,展开为单个记录的ID+日期 valid_records = df_with_diff.filter((F.col("days_diff") < 15) & (F.col("sum_diff") <= 3)) \ .select("ID", "ContractDate", "prev_date") \ .withColumn("valid_date", F.explode(F.array("ContractDate", "prev_date"))) \ .select("ID", "valid_date") \ .distinct()
4. 关联原表获取最终结果
将原DataFrame与有效记录标识关联,保留匹配的记录:
# 关联并选择原表所有列 result_df = df.join( valid_records, on=[df.ID == valid_records.ID, df.ContractDate == valid_records.valid_date], how="inner" ).select(df["*"]) # 展示结果 result_df.show()
执行后即可得到你期望的结果:
+----+---+-------------+-------+-------+ |Name| ID|ContractDate |LoanSum| Status| +----+---+-------------+-------+-------+ | A|ID1| 2022-10-10| 10| Closed| | A|ID1| 2022-10-15| 13| Active| | C|ID3| 2022-12-10| 40| Closed| | C|ID3| 2022-12-12| 43| Active| | C|ID3| 2022-12-19| 46| Active| +----+---+-------------+-------+-------+
内容的提问来源于stack exchange,提问作者lenpyspanacb
相关产品推荐
相关产品推荐

