简介:这份资源是面向大数据与数据分析初学者的课程设计实战包,以和鲸社区信用卡评分模型数据为数据集,用Python结合Spark完成数据预处理、指标分析与结果可视化,适合正在学习Spark框架、需要完整项目案例练手的高校学生和开发者参考。压缩包共22个文件,约4.91MB,包含4个Python脚本分别负责数据分析、数据预处理与Web可视化,2个CSV数据文件提供原始与清洗后数据,另有5个HTML可视化结果页面、5个XML配置文件和1份课程设计报告文档,结构完整、层次清晰。目前已有3779人学习下载,说明该案例具备一定的参考价值。读者可以从中获得一套可直接运行的Spark分析流程、数据清洗与统计思路、可视化页面生成方法,以及一份可对照的课程设计报告,便于快速理解项目整体架构并迁移到自己的分析任务中。
1. 信用卡评分卡遇上 Spark:一份能跑通全流程的数据分析资源
信用卡评分这件事,业务侧关心的是「这个人会不会逾期」,技术侧关心的是「几十万条申请记录、上百万条还款流水,怎么在可接受的时间里跑出稳定的分数」。单机 pandas 处理几十万行就开始喘,特征工程一上 groupby 加窗口函数,内存直接爆掉,这是很多人从「会数据分析」到「能做评分卡」之间翻车的地方。这份基于 Spark 的信用卡评分数据分析资源,解决的正是这个断层:它把数据清洗、特征衍生、WOE/IV 计算、逻辑回归建模、评分映射这一整条链路,用 Spark 的 DataFrame 和 MLlib 串了起来,适合已经会 Python、想往大数据风控方向落地的从业者,也适合拿它当 Spark 数据分析项目练手的人。下面我按「资源是什么、怎么用、坑在哪」拆开讲。
2. 环境与数据准备:从 SparkSession 到评分卡原始表
2.1 为什么评分卡项目要用 Spark 而不是 pandas
先讲选型理由,不然很多人跑一半会怀疑自己是不是过度设计。信用卡评分的数据形态通常是三张表:申请信息表(客户基本属性、申请额度)、还款行为表(每月账单、还款状态、逾期天数)、交易流水表(消费金额、商户类型)。行为表和流水表的量级往往是申请表的几十倍,做「近 6 个月最大逾期天数」「近 3 个月平均使用额度」这类特征时,需要对每个客户做时间窗口聚合。pandas 的做法是先 groupby 再 rolling,数据量一大就吃满内存;Spark 的窗口函数和 DataFrame API 天然按分区并行,同样的逻辑可以横向扩机器。
另一个理由是特征衍生阶段的可复现性。评分卡对特征口径极其敏感,同一个「逾期次数」定义差一天,IV 值就变了。Spark 的 SQL 表达和 DataFrame 转换可以固化成脚本,配合版本管理,比在 notebook 里手写 pandas 更容易复现。常见做法是把特征逻辑写成一段 Spark SQL,参数(观察期、表现期)抽成变量,换一批数据只改参数不改逻辑。
2.2 环境搭建与依赖确认
这份资源跑在本地或单机伪分布式都能起步,不强制上集群。核心依赖是 PySpark,版本建议 3.x,Python 3.8 以上。先确认环境:
# 检查 Java,Spark 依赖 JVM java -version # 检查 Python python --version # 安装 PySpark,指定版本避免和集群不一致 pip install pyspark==3.5.0Java 版本要注意,Spark 3.x 配 JDK 8 或 11 都行,但 JDK 17 在部分老版本上会有模块访问报错,这是血泪经验,别一上来就用最新 JDK。装完写个最小验证:
from pyspark.sql import SparkSession # 本地模式启动,master 用 local[*] 吃满本机核数 spark = SparkSession.builder \ .appName("credit_scorecard") \ .master("local[*]") \ .config("spark.sql.shuffle.partitions", "8") \ .getOrCreate() # 打印版本,确认环境 print(spark.version)spark.sql.shuffle.partitions默认是 200,本地跑小数据时 200 个分区会让每个任务只处理几行,调度开销反而拖慢速度,设成核数的 2 到 4 倍比较合适。这一步很多人忽略,然后抱怨「Spark 比 pandas 还慢」,其实是分区数没调。
2.3 数据加载与字段规整
评分卡原始数据常见两种来源:CSV 落地文件,或者 Hive 表。资源里给的是 CSV 加载路径,读进来先做类型规整,因为金额、天数这类字段经常被当成字符串。
from pyspark.sql import functions as F from pyspark.sql.types import IntegerType, DoubleType # 读取申请与行为数据,inferSchema 会多扫一遍,生产上建议显式指定 schema app_df = spark.read.csv("data/application.csv", header=True, inferSchema=True) behavior_df = spark.read.csv("data/behavior.csv", header=True, inferSchema=True) # 统一逾期天数为整型,金额为浮点,空值先置 0 再判断 behavior_df = behavior_df \ .withColumn("overdue_days", F.col("overdue_days").cast(IntegerType())) \ .withColumn("bill_amount", F.col("bill_amount").cast(DoubleType())) \ .fillna({"overdue_days": 0, "bill_amount": 0.0}) # 按客户 ID 关联,注意一对多关系,行为表一个客户多条 joined = app_df.join(behavior_df, on="cust_id", how="left") print(joined.count())inferSchema=True在开发阶段方便,但生产上会多一次全表扫描,数据量大时建议手写 StructType。fillna的顺序有讲究:先把 null 填成 0 再参与后续聚合,否则sum遇到 null 结果还是 null,特征直接缺失。join 用 left 保留所有申请客户,行为缺失的客户在评分卡里通常按「无行为」单独分箱,不能直接丢,丢了会引入样本偏差。
3. 特征工程:用 Spark SQL 做时间窗口聚合与 WOE 计算
3.1 时间窗口特征衍生
评分卡的核心特征几乎都带时间窗口:近 6 个月、近 12 个月、近 24 个月。用 Spark 的窗口函数可以一次算出多个口径。下面这段是「近 6 个月最大逾期天数」和「近 3 个月平均账单」的典型写法:
from pyspark.sql import Window # 按客户分区,按账期倒序,取最近 N 期 w = Window.partitionBy("cust_id").orderBy(F.col("bill_month").desc()) feat_df = behavior_df \ .withColumn("rn", F.row_number().over(w)) \ .withColumn("max_overdue_6m", F.when(F.col("rn") <= 6, F.col("overdue_days")).otherwise(None)) \ .withColumn("avg_bill_3m", F.when(F.col("rn") <= 3, F.col("bill_amount")).otherwise(None)) \ .groupBy("cust_id") \ .agg( F.max("max_overdue_6m").alias("max_overdue_6m"), F.avg("avg_bill_3m").alias("avg_bill_3m") )逻辑说明:先用row_number给每个客户的账单按月倒序编号,编号 1 是最近一期。然后用when把窗口外的行置为 null,这样聚合时max和avg自动忽略它们。参数上,rn <= 6就是观察期长度,改成 12 就是近 12 个月,换口径只改这个数字。注意max对 null 的处理是忽略,但如果一个客户 6 期内全是 null,结果就是 null,后面要单独处理。
3.2 WOE 与 IV 的计算
WOE(Weight of Evidence)和 IV(Information Value)是评分卡特征筛选的标准工具。Spark 里没有现成函数,得自己按分箱算。核心是统计每个箱内的好样本和坏样本数:
# 假设 label 列:1 为坏客户(逾期),0 为好客户 def calc_woe_iv(df, feature, label, bins=10): # 等频分箱,用 approxQuantile 拿分位点 splits = df.approxQuantile(feature, [i / bins for i in range(1, bins)], 0.01) # 用 Bucketizer 分箱 from pyspark.ml.feature import Bucketizer bucketizer = Bucketizer(splits=[float("-inf")] + splits + [float("inf")], inputCol=feature, outputCol=feature + "_bin") binned = bucketizer.transform(df) # 统计每箱好坏样本 stat = binned.groupBy(feature + "_bin").agg( F.sum(F.when(F.col(label) == 1, 1).otherwise(0)).alias("bad"), F.sum(F.when(F.col(label) == 0, 1).otherwise(0)).alias("good") ).collect() total_bad = sum(r["bad"] for r in stat) total_good = sum(r["good"] for r in stat) iv = 0.0 for r in stat: bad_rate = (r["bad"] + 0.5) / total_bad good_rate = (r["good"] + 0.5) / total_good woe = math.log(good_rate / bad_rate) iv += (good_rate - bad_rate) * woe return ivapproxQuantile的第三个参数是相对误差,0.01 表示分位点允许 1% 误差,数据量大时调大能提速。分子分母加 0.5 是拉普拉斯平滑,防止某个箱好或坏样本为 0 导致 log 无定义,这是评分卡里的常规操作。IV 值一般大于 0.02 才考虑保留,大于 0.5 要警惕,可能是特征穿越了未来信息。
3.3 特征筛选与共线性处理
算出 IV 后不是全塞进模型。常见做法是:先按 IV 排序,保留 IV 在 0.02 到 0.5 之间的;再算两两相关系数,相关系数超过 0.7 的保留 IV 高的那个。Spark 的Correlation.corr可以算皮尔逊相关:
from pyspark.ml.stat import Correlation from pyspark.ml.feature import VectorAssembler selected = ["max_overdue_6m", "avg_bill_3m", "credit_limit", "age"] assembler = VectorAssembler(inputCols=selected, outputCol="features") vec_df = assembler.transform(feat_df).select("features") corr_matrix = Correlation.corr(vec_df, "features").head()[0] print(corr_matrix.toArray())这一步的意义在于,逻辑回归对共线性敏感,两个高度相关的特征会让系数不稳定,换一批样本系数符号都可能翻。筛完特征再进模型,比一股脑全丢进去稳得多。
4. 建模与评分映射:MLlib 逻辑回归到标准分
4.1 训练集与测试集划分
评分卡建模要按时间切,不能随机切。随机切会让未来样本混进训练集,评估结果虚高。正确做法是按申请月份切,比如前 18 个月训练,后 6 个月测试:
train = feat_df.filter(F.col("apply_month") <= "2023-06") test = feat_df.filter(F.col("apply_month") > "2023-06") print("train:", train.count(), "test:", test.count())如果样本坏客户占比很低(比如 1%),还要做样本加权,给坏客户更高权重,否则模型会倾向于全预测为好客户。MLlib 的LogisticRegression支持weightCol参数。
4.2 逻辑回归训练与参数
from pyspark.ml.classification import LogisticRegression from pyspark.ml.feature import VectorAssembler assembler = VectorAssembler(inputCols=selected, outputCol="features") train_vec = assembler.transform(train).select("features", "label") test_vec = assembler.transform(test).select("features", "label") lr = LogisticRegression( featuresCol="features", labelCol="label", maxIter=100, regParam=0.01, # L2 正则,防过拟合 elasticNetParam=0.0 # 0 为纯 L2 ) model = lr.fit(train_vec)regParam是正则强度,评分卡样本量通常几万到几十万,0.01 到 0.1 之间试。maxIter设 100 一般够收敛,如果日志里看到「not converged」,先加迭代次数,再检查特征是不是没标准化。逻辑回归对量纲敏感,金额和年龄放一起,梯度下降会震荡,建议先做标准化或用StandardScaler。
4.3 评分映射:把概率转成 300 到 850 的分数
业务不认概率,认分数。标准做法是设定基准分和 PDO(Points to Double the Odds):
import math base_score = 600 # 基准分 base_odds = 50 # 基准分对应的好坏比 pdo = 20 # 好坏比翻倍需要的分数 factor = pdo / math.log(2) offset = base_score - factor * math.log(base_odds) def prob_to_score(p): odds = (1 - p) / p return offset + factor * math.log(odds)factor和offset是评分卡的标定参数,PDO 越小分数对风险越敏感。算完分数后,通常还要按分数段做通过率、坏账率的回溯,确认单调性:分数越高,坏账率越低,如果出现倒挂,说明特征或分箱有问题。
5. 避坑与排查:评分卡跑 Spark 最常见的五个翻车点
5.1 现象:任务卡在某个 stage 不动,日志刷 shuffle
原因:join 或 groupBy 触发了大量 shuffle,分区键倾斜,某个 key 的数据量远超其他。信用卡数据里,如果按商户 ID 聚合,头部商户可能占一半流水。
解决:先看 Spark UI 里各 task 的耗时分布,明显偏长的就是倾斜分区。对倾斜 key 加随机前缀打散,聚合两次;或者用repartition按更均匀的列重分区。别一上来就加内存,倾斜是分布问题不是资源问题。
5.2 现象:WOE 计算报 math domain error
原因:某个分箱里好样本或坏样本为 0,log 里出现 0 或负数。
解决:分子分母加平滑项(前面代码里的 0.5),或者把样本过少的箱合并到相邻箱。分箱数别设太多,10 箱在几万样本下已经够细,箱里样本太少统计不稳定。
5.3 现象:训练集 AUC 0.85,测试集掉到 0.6
原因:特征穿越。比如用了「当前逾期状态」去预测「是否逾期」,或者时间窗口没对齐,把表现期的信息混进了观察期。
解决:严格按时间切分观察期和表现期,特征只能用观察期结束前的数据。每个特征过一遍「这个值在申请时点能不能拿到」,拿不到的一律删。
5.4 现象:本地跑得好,上集群报序列化错误
原因:UDF 里引用了不可序列化的对象,比如数据库连接、文件句柄。
解决:UDF 里只做纯计算,外部资源在 driver 侧准备好再广播。能用内置函数就别写 UDF,UDF 会破坏 Catalyst 优化,性能差一截。
5.5 现象:分数分布全挤在中间,区分度差
原因:特征没标准化,或者正则太强把系数压得太小。
解决:检查StandardScaler是否加在正确位置,regParam调小试。另外确认标签定义是否清晰,坏客户口径模糊会让模型学不到东西。
6. 进阶技巧:用分箱单调性校验和 PSI 监控守住线上分数
模型上线不是终点。评分卡最怕的是分数漂移,今天批的客户和三个月前分布不一样,模型还在按老规律打分。这里给两个我每次上线都强制走的检查。
第一个是分箱单调性校验。算完 WOE 后,把每个特征的箱按 WOE 排序,看坏样本率是否单调。如果出现「中间箱坏账率反而低」的倒挂,说明分箱不合理,常见于把缺失值单独分了一箱但没处理好。下面这段可以快速输出每个箱的坏账率:
# 假设 binned 是分箱后的 DataFrame,label 为 1 是坏 check = binned.groupBy(feature + "_bin").agg( F.count("*").alias("cnt"), F.mean("label").alias("bad_rate") ).orderBy(feature + "_bin") check.show()看bad_rate那一列,正常应该随箱号递增或递减,出现 V 形就要回去查分箱边界。这个检查花不了几分钟,但能挡掉很多「模型指标好看、上线就翻车」的情况。
第二个是 PSI(Population Stability Index)监控。上线后每月拿新样本的分数分布和训练集比,PSI 小于 0.1 算稳定,0.1 到 0.25 要警惕,超过 0.25 就得考虑重新训练。计算方式是按分数分箱,比较两期的占比:
def calc_psi(base_df, curr_df, score_col, bins=10): # 用训练集分位点做统一分箱边界 splits = base_df.approxQuantile(score_col, [i / bins for i in range(1, bins)], 0.01) from pyspark.ml.feature import Bucketizer b = Bucketizer(splits=[float("-inf")] + splits + [float("inf")], inputCol=score_col, outputCol="bin") base_pct = b.transform(base_df).groupBy("bin").count().toPandas() curr_pct = b.transform(curr_df).groupBy("bin").count().toPandas() base_pct["pct"] = base_pct["count"] / base_pct["count"].sum() curr_pct["pct"] = curr_pct["count"] / curr_pct["count"].sum() merged = base_pct.merge(curr_pct, on="bin", suffixes=("_base", "_curr")) psi = ((merged["pct_curr"] - merged["pct_base"]) * np.log(merged["pct_curr"] / merged["pct_base"])).sum() return psi分箱边界必须用训练集的,不能用当前期的,否则两期边界不一致,PSI 没意义。这个函数我一般挂到调度里,每月自动跑一次,PSI 超阈值就告警。
还有一个容易被忽略的点:评分卡的特征口径要写进文档,和业务对齐。技术侧觉得「逾期天数」很明确,业务侧可能理解成「当前逾期」还是「历史最大逾期」,差一个字结果完全不同。我现在的习惯是每个特征上线前,拉业务方确认一遍口径,确认完再固化到脚本里,从那以后每次改特征都强制走一遍这个确认,省掉了很多事后扯皮。希望这份资源能帮你把评分卡这条链路真正跑通,而不只是停在 notebook 里。
本文还有配套的精品资源,点击获取