求助:使用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出错,大概率是这几个原因:
- 关联键类型不匹配:贷款表的
PinCode是STRING,征信表是LONG,直接关联会导致完全匹配不上; - 关联键选得不对:如果用
CustomerName或者City这种非唯一字段,会出现错误关联; - 过滤逻辑顺序错了:应该先过滤征信表的符合条件数据再关联,既减少计算量,也避免关联后过滤的逻辑混乱;
- 重复列没处理:两张表有
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
相关产品推荐
相关产品推荐

