简介这份资源是面向大数据与数据分析初学者的课程设计实战包以和鲸社区信用卡评分模型数据为数据集用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.xPython 3.8 以上。先确认环境# 检查 JavaSpark 依赖 JVM java -version # 检查 Python python --version # 安装 PySpark指定版本避免和集群不一致 pip install pyspark3.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, headerTrue, inferSchemaTrue) behavior_df spark.read.csv(data/behavior.csv, headerTrue, inferSchemaTrue) # 统一逾期天数为整型金额为浮点空值先置 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, oncust_id, howleft) print(joined.count())inferSchemaTrue在开发阶段方便但生产上会多一次全表扫描数据量大时建议手写 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 的计算WOEWeight of Evidence和 IVInformation Value是评分卡特征筛选的标准工具。Spark 里没有现成函数得自己按分箱算。核心是统计每个箱内的好样本和坏样本数# 假设 label 列1 为坏客户逾期0 为好客户 def calc_woe_iv(df, feature, label, bins10): # 等频分箱用 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)], inputColfeature, outputColfeature _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(inputColsselected, outputColfeatures) 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(inputColsselected, outputColfeatures) train_vec assembler.transform(train).select(features, label) test_vec assembler.transform(test).select(features, label) lr LogisticRegression( featuresColfeatures, labelCollabel, maxIter100, regParam0.01, # L2 正则防过拟合 elasticNetParam0.0 # 0 为纯 L2 ) model lr.fit(train_vec)regParam是正则强度评分卡样本量通常几万到几十万0.01 到 0.1 之间试。maxIter设 100 一般够收敛如果日志里看到「not converged」先加迭代次数再检查特征是不是没标准化。逻辑回归对量纲敏感金额和年龄放一起梯度下降会震荡建议先做标准化或用StandardScaler。4.3 评分映射把概率转成 300 到 850 的分数业务不认概率认分数。标准做法是设定基准分和 PDOPoints to Double the Oddsimport 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原因某个分箱里好样本或坏样本为 0log 里出现 0 或负数。解决分子分母加平滑项前面代码里的 0.5或者把样本过少的箱合并到相邻箱。分箱数别设太多10 箱在几万样本下已经够细箱里样本太少统计不稳定。5.3 现象训练集 AUC 0.85测试集掉到 0.6原因特征穿越。比如用了「当前逾期状态」去预测「是否逾期」或者时间窗口没对齐把表现期的信息混进了观察期。解决严格按时间切分观察期和表现期特征只能用观察期结束前的数据。每个特征过一遍「这个值在申请时点能不能拿到」拿不到的一律删。5.4 现象本地跑得好上集群报序列化错误原因UDF 里引用了不可序列化的对象比如数据库连接、文件句柄。解决UDF 里只做纯计算外部资源在 driver 侧准备好再广播。能用内置函数就别写 UDFUDF 会破坏 Catalyst 优化性能差一截。5.5 现象分数分布全挤在中间区分度差原因特征没标准化或者正则太强把系数压得太小。解决检查StandardScaler是否加在正确位置regParam调小试。另外确认标签定义是否清晰坏客户口径模糊会让模型学不到东西。6. 进阶技巧用分箱单调性校验和 PSI 监控守住线上分数模型上线不是终点。评分卡最怕的是分数漂移今天批的客户和三个月前分布不一样模型还在按老规律打分。这里给两个我每次上线都强制走的检查。第一个是分箱单调性校验。算完 WOE 后把每个特征的箱按 WOE 排序看坏样本率是否单调。如果出现「中间箱坏账率反而低」的倒挂说明分箱不合理常见于把缺失值单独分了一箱但没处理好。下面这段可以快速输出每个箱的坏账率# 假设 binned 是分箱后的 DataFramelabel 为 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 形就要回去查分箱边界。这个检查花不了几分钟但能挡掉很多「模型指标好看、上线就翻车」的情况。第二个是 PSIPopulation Stability Index监控。上线后每月拿新样本的分数分布和训练集比PSI 小于 0.1 算稳定0.1 到 0.25 要警惕超过 0.25 就得考虑重新训练。计算方式是按分数分箱比较两期的占比def calc_psi(base_df, curr_df, score_col, bins10): # 用训练集分位点做统一分箱边界 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)], inputColscore_col, outputColbin) 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, onbin, suffixes(_base, _curr)) psi ((merged[pct_curr] - merged[pct_base]) * np.log(merged[pct_curr] / merged[pct_base])).sum() return psi分箱边界必须用训练集的不能用当前期的否则两期边界不一致PSI 没意义。这个函数我一般挂到调度里每月自动跑一次PSI 超阈值就告警。还有一个容易被忽略的点评分卡的特征口径要写进文档和业务对齐。技术侧觉得「逾期天数」很明确业务侧可能理解成「当前逾期」还是「历史最大逾期」差一个字结果完全不同。我现在的习惯是每个特征上线前拉业务方确认一遍口径确认完再固化到脚本里从那以后每次改特征都强制走一遍这个确认省掉了很多事后扯皮。希望这份资源能帮你把评分卡这条链路真正跑通而不只是停在 notebook 里。本文还有配套的精品资源点击获取