PySpark多DataFrame按条件匹配:添加指定规则的PaymentSum列
PySpark解决方案:匹配最接近到期日的还款金额
注意:从你的目标结果来看,实际需求应为还款日期最接近贷款的到期日(MaturityDate),而非合同日期(ContractDate),以下解决方案将基于此逻辑实现。
步骤说明
- 转换日期类型:将字符串格式的日期转为PySpark的
DateType,支持日期计算。 - 关联数据表:通过贷款ID关联贷款基本信息表和还款记录表。
- 计算日期差:对非D银行的贷款,计算还款日期与到期日的天数差绝对值。
- 筛选最优记录:按贷款ID分组,筛选出日期差最小的还款记录;若存在多条相同差的记录,取最晚还款的那条(匹配你的目标结果)。
- 合并结果:将筛选出的还款金额合并回原始贷款表,D银行的贷款该列留空。
完整代码实现
from pyspark.sql import SparkSession from pyspark.sql.functions import col, datediff, abs, row_number from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("LoanPaymentMatch").getOrCreate() # 1. 创建示例贷款基本信息DataFrame loan_data = [ ("ID1", "2024-06-01", "2024-06-18", "A"), ("ID2", "2024-06-05", "2024-06-18", "B"), ("ID3", "2024-06-10", "2024-06-17", "C"), ("ID4", "2024-06-15", None, "D"), ("ID5", "2024-08-01", "2024-08-22", "A"), ("ID6", "2024-08-08", "2024-08-23", "B"), ("ID7", "2024-08-20", None, "D") ] loan_df = spark.createDataFrame(loan_data, ["ID", "ContractDate", "MaturityDate", "Bank"]) # 转换日期列类型 loan_df = loan_df.withColumn("ContractDate", col("ContractDate").cast("date")) \ .withColumn("MaturityDate", col("MaturityDate").cast("date")) # 2. 创建示例还款记录DataFrame payment_data = [ ("ID1", "2024-06-02", 10), ("ID1", "2024-06-08", 40), ("ID1", "2024-06-10", 50), ("ID2", "2024-06-06", 30), ("ID2", "2024-06-07", 90), ("ID2", "2024-06-08", 20), ("ID3", "2024-06-11", 20), ("ID3", "2024-06-12", 30), ("ID3", "2024-06-13", 50), ("ID5", "2024-08-10", 15), ("ID5", "2024-08-13", 35), ("ID5", "2024-08-15", 30), ("ID6", "2024-08-15", 20), ("ID6", "2024-08-16", 20), ("ID6", "2024-08-20", 70) ] payment_df = spark.createDataFrame(payment_data, ["ID_loan", "PaymentDate", "PaymentSum"]) # 转换还款日期类型 payment_df = payment_df.withColumn("PaymentDate", col("PaymentDate").cast("date")) # 3. 关联两张表,计算日期差 joined_df = loan_df.join(payment_df, loan_df.ID == payment_df.ID_loan, "left") \ .withColumn("days_diff", abs(datediff(col("MaturityDate"), col("PaymentDate")))) \ .filter(col("Bank") != "D") # 只处理非D银行的贷款 # 4. 定义窗口函数,筛选每个贷款的最优还款记录 window_spec = Window.partitionBy("ID").orderBy(col("days_diff").asc(), col("PaymentDate").desc()) ranked_df = joined_df.withColumn("rn", row_number().over(window_spec)) \ .filter(col("rn") == 1) \ .select("ID", "PaymentSum") # 5. 合并回原始贷款表,处理D银行的情况 result_df = loan_df.join(ranked_df, on="ID", how="left") \ .select("ID", "ContractDate", "MaturityDate", "Bank", "PaymentSum") # 展示结果 result_df.show()
代码解释
- 日期转换:通过
cast("date")将字符串日期转为日期类型,确保后续datediff函数能正常计算天数差。 - 窗口函数:
partitionBy("ID")按贷款ID分组,orderBy(col("days_diff").asc(), col("PaymentDate").desc())先按日期差升序(取最小差),再按还款日期降序(差相同时取最晚还款的记录)。 - 关联合并:最后通过左关联将筛选出的还款金额合并回原始贷款表,D银行的贷款因未参与前面的计算,
PaymentSum列自动为null,符合需求。
内容的提问来源于stack exchange,提问作者lenpyspanacb
相关产品推荐
相关产品推荐

