PySpark高级过滤:计算符合特定条件的借款人贷款总额
解决方案
要实现你的需求,需针对同一借款人+同一银行的Active状态贷款,筛选出首笔发放日起3天内累计达3笔及以上的组,再计算这些组的LoanSum总和。以下是完整实现步骤:
步骤1:过滤Active状态的贷款
先筛选出状态为Active的记录,排除已结清的贷款:
from pyspark.sql import functions as f from pyspark.sql import Window # 假设data是你的原始数据列表 df = spark.createDataFrame(data).toDF('ID','ContractDate','LoanSum','ClosingDate', 'Status', 'Bank') # 过滤Active状态的记录 df_active = df.filter(f.col("Status") == "Active")
步骤2:计算每个ID+Bank组的首笔贷款日期
按ID和Bank分区(确保同一借款人同一银行),按ContractDate排序后,提取每组的首笔贷款日期:
# 定义窗口:按ID+Bank分区,按ContractDate升序排序 w_id_bank = Window.partitionBy("ID", "Bank").orderBy("ContractDate") # 添加首笔贷款日期列 df_with_first_date = df_active.withColumn( "first_contract_date", f.first("ContractDate").over(w_id_bank) # 取每组最早的ContractDate )
步骤3:筛选首笔3天内的贷款
计算每条记录与首笔贷款的日期差,保留差≤3天的记录:
df_3days_window = df_with_first_date.filter( f.datediff(f.col("ContractDate"), f.col("first_contract_date")) <= 3 )
步骤4:统计符合条件的组并求和
按ID+Bank分组,统计3天内的贷款笔数,筛选笔数≥3的组,同时计算LoanSum总和:
result = df_3days_window.groupBy("ID", "Bank").agg( f.count("*").alias("loan_count"), f.sum("LoanSum").alias("total_loan_sum") ).filter(f.col("loan_count") >= 3) result.show()
结果验证
运行上述代码后,结果会输出:
+---+----+-----------+---------------+ | ID|Bank|loan_count|total_loan_sum| +---+----+-----------+---------------+ |ID3| A| 3| 100| +---+----+-----------+---------------+
完全匹配你示例中的预期结果。
内容的提问来源于stack exchange,提问作者lenpyspanacb
相关产品推荐
相关产品推荐

