PySpark复杂过滤操作求助:借款人贷款数据多条件筛选
PySpark复杂多条件筛选实现
原始数据准备
首先创建测试用的DataFrame,模拟用户提供的贷款数据:
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, IntegerType, DateType from pyspark.sql import functions as F spark = SparkSession.builder.appName("LoanFilter").getOrCreate() schema = StructType([ StructField("ID", StringType(), True), StructField("ContractDate", StringType(), True), StructField("Loansum", IntegerType(), True), StructField("Debt", IntegerType(), True), StructField("MaturityDate", StringType(), True), StructField("Bank", StringType(), True) ]) data = [ ("ID1", "2024-06-01", 100, 10, "2024-06-18", "A"), ("ID1", "2024-06-05", 50, 20, "2024-06-17", "B"), ("ID1", "2024-06-10", 50, 20, "2024-06-15", "C"), ("ID1", "2024-06-10", 90, 70, "2024-06-30", "D"), ("ID1", "2024-06-15", 50, 50, None, "R"), ("ID1", "2024-07-15", 100, 90, None, "D"), ("ID2", "2024-08-01", 70, 20, "2024-09-01", "A"), ("ID2", "2024-08-08", 50, 10, "2024-08-26", "B"), ("ID2", "2024-08-20", 32, 32, None, "R"), ("ID3", "2024-09-01", 60, 10, "2024-09-24", "A"), ("ID3", "2024-09-03", 40, 10, "2024-09-23", "B"), ("ID3", "2024-09-22", 22, 22, None, "R"), ("ID4", "2024-10-03", 40, 10, "2024-10-23", "B"), ("ID5", "2024-11-01", 70, 20, "2024-11-21", "A"), ("ID5", "2024-11-08", 50, 10, "2024-11-22", "B"), ("ID5", "2024-11-20", 40, 40, None, "R") ] df = spark.createDataFrame(data, schema=schema)
步骤1:日期类型转换
将字符串格式的日期字段转为DateType,方便后续日期计算:
df = df.withColumn("ContractDate", F.to_date(F.col("ContractDate"), "yyyy-MM-dd")) \ .withColumn("MaturityDate", F.to_date(F.col("MaturityDate"), "yyyy-MM-dd"))
步骤2:提取R银行贷款信息
筛选出所有银行R的贷款记录,作为后续条件判断的基准:
r_loans = df.filter(F.col("Bank") == "R") \ .select("ID", "ContractDate", "Loansum") \ .withColumnRenamed("ContractDate", "R_ContractDate") \ .withColumnRenamed("Loansum", "R_Loansum")
步骤3:筛选符合条件的前置贷款
关联主表与R贷款表,筛选出R贷款发放前的记录,且到期日不晚于R贷款合同日期后5天:
joined_df = df.join(r_loans, on="ID", how="left") \ .filter(F.col("ContractDate") < F.col("R_ContractDate")) valid_previous = joined_df.filter( F.col("MaturityDate") <= F.date_add(F.col("R_ContractDate"), 5) )
步骤4:验证R贷款金额条件
按ID分组计算符合条件的前置贷款Debt总和,判断R贷款金额是否不超过该总和的110%:
sum_debt_df = valid_previous.groupBy("ID") \ .agg(F.sum("Debt").alias("TotalValidDebt")) valid_ids = sum_debt_df.join(r_loans, on="ID") \ .filter(F.col("R_Loansum") <= F.col("TotalValidDebt") * 1.1) \ .select("ID")
步骤5:标记并合并结果
将符合条件的前置贷款标记为Previous,R贷款标记为Last,合并后排序输出:
previous_final = valid_previous.join(valid_ids, on="ID") \ .select("ID", "ContractDate", "Loansum", "Debt", "MaturityDate", "Bank") \ .withColumn("Marker", F.lit("Previous")) r_final = r_loans.join(valid_ids, on="ID") \ .join(df, on=["ID", "ContractDate"], how="left") \ .select("ID", "ContractDate", "Loansum", "Debt", "MaturityDate", "Bank") \ .withColumn("Marker", F.lit("Last")) final_df = previous_final.union(r_final).orderBy("ID", "ContractDate") final_df.show(truncate=False)
执行上述代码后,输出结果将与预期一致。
内容的提问来源于stack exchange,提问作者lenpyspanacb
相关产品推荐
相关产品推荐

