You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

求助:使用PySpark Core和SQL实现贷款审批状态判定遇问题

PySpark Core & SQL 实现贷款审批状态判定

问题背景

我现在要处理两张客户数据表,需求是:当客户非违约(DefaulterFlag='N')且CibilScore>800时,标记贷款状态为Approved。用PySpark Core实现时得到错误结果,但导入MySQL用SQL跑是对的,求正确的PySpark实现方案。

两张表的结构和示例数据如下:

  • 客户贷款信息表:
    Schema:CustomerName:STRING, DOB:STRING, UIN:STRING, MailID:STRING, PhoneNumber:LONG, City:STRING, State:STRING, LivingStatus:STRING, PinCode:STRING, LoanAmount:LONG
    示例数据:
    Sakshi, 22-03-86, UIN0043, Sakshi@mail.com, 3344990876, Ahmedabad, Gujarat, BPL ,380001, 23000
    Shivani, 22-02-83, UIN0044, Shivani@mail.com, 3344990876, Thiruvananthpuram, Kerala, APL, 695001,24500
    
  • 客户征信信息表:
    Schema:CustomerName:STRING, DOB:STRING,UIN:STRING, City:STRING, State:STRING, PinCode:LONG, CibilScore:LONG, DefaulterFlag:STRING
    示例数据:
    Shubham, 23-08-86, UIN0007, Thiruvananthpuram, Kerala, 695001, 3530, N
    Anushka, 25-08-82, UIN0008, Thiruvananthpuram, Kerala, 695001, 1530, Y
    

你可能踩的坑

PySpark Core出错,大概率是这几个原因:

  1. 关联键类型不匹配:贷款表的PinCode是STRING,征信表是LONG,直接关联会导致完全匹配不上;
  2. 关联键选得不对:如果用CustomerName或者City这种非唯一字段,会出现错误关联;
  3. 过滤逻辑顺序错了:应该先过滤征信表的符合条件数据再关联,既减少计算量,也避免关联后过滤的逻辑混乱;
  4. 重复列没处理:两张表有CustomerName、DOB这类重复列,关联后容易搞混字段。

正确实现代码

先初始化SparkSession

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, when

spark = SparkSession.builder.appName("LoanApprovalCheck").getOrCreate()

加载数据(以CSV为例)

# 加载贷款信息表,指定schema避免自动推断出错
loan_df = spark.read.csv(
    "loan_data.csv", 
    header=False, 
    schema="CustomerName:STRING, DOB:STRING, UIN:STRING, MailID:STRING, PhoneNumber:LONG, City:STRING, State:STRING, LivingStatus:STRING, PinCode:STRING, LoanAmount:LONG"
)

# 加载征信信息表
cibil_df = spark.read.csv(
    "cibil_data.csv", 
    header=False, 
    schema="CustomerName:STRING, DOB:STRING,UIN:STRING, City:STRING, State:STRING, PinCode:LONG, CibilScore:LONG, DefaulterFlag:STRING"
)

方法一:PySpark Core(RDD API)实现

# 1. 转换RDD,提取需要的字段,同时把征信表的PinCode转成STRING,统一类型
loan_rdd = loan_df.rdd.map(lambda row: (row.UIN, (row.CustomerName, row.DOB, row.LoanAmount, row.PinCode)))
cibil_rdd = cibil_df.rdd.map(lambda row: (row.UIN, (row.CibilScore, row.DefaulterFlag, str(row.PinCode))))

# 2. 用UIN关联(UIN是客户唯一标识,比其他字段可靠)
joined_rdd = loan_rdd.join(cibil_rdd)

# 3. 过滤符合条件的记录,添加贷款状态
approved_rdd = joined_rdd.filter(lambda x: x[1][1][1] == 'N' and x[1][1][0] > 800)\
                         .map(lambda x: (
                             x[1][0][0],  # 客户姓名
                             x[1][0][1],  # 出生日期
                             x[0],        # UIN
                             x[1][0][2],  # 贷款金额
                             x[1][1][0],  # 征信分数
                             x[1][1][1],  # 是否违约
                             "Approved"   # 贷款状态
                         ))

# 4. 转成DataFrame展示结果
result_core_df = approved_rdd.toDF(["CustomerName", "DOB", "UIN", "LoanAmount", "CibilScore", "DefaulterFlag", "LoanStatus"])
result_core_df.show()

方法二:PySpark SQL实现

# 1. 注册临时视图,方便写SQL查询
loan_df.createOrReplaceTempView("loan_table")
cibil_df.createOrReplaceTempView("cibil_table")

# 2. 编写SQL,注意统一PinCode类型,用UIN关联,添加条件判断
sql_query = """
SELECT 
    l.CustomerName,
    l.DOB,
    l.UIN,
    l.LoanAmount,
    c.CibilScore,
    c.DefaulterFlag,
    CASE 
        WHEN c.DefaulterFlag = 'N' AND c.CibilScore > 800 THEN 'Approved'
        ELSE 'Rejected'
    END AS LoanStatus
FROM loan_table l
INNER JOIN cibil_table c 
    ON l.UIN = c.UIN 
    AND l.PinCode = CAST(c.PinCode AS STRING)
"""

# 3. 执行查询并展示结果
result_sql_df = spark.sql(sql_query)
result_sql_df.show()

额外优化建议

  • 提前过滤数据:在Core API里,可以先过滤征信表中符合条件的记录再关联,减少计算量:
    filtered_cibil_rdd = cibil_rdd.filter(lambda x: x[1][1] == 'N' and x[1][0] > 800)
    joined_rdd = loan_rdd.join(filtered_cibil_rdd)
    
  • 避免笛卡尔积:如果关联键选得不对,很容易出现大量错误匹配,一定要用唯一标识(比如UIN)来关联。

内容的提问来源于stack exchange,提问作者EdwardFunnyHands

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.23 05:02:11