
简介这份资源是面向大数据与数据分析初学者、高校课程设计参考者的Spark实战项目以和鲸社区信用卡评分模型构建数据为数据集用Python结合Spark完成数据预处理、统计分析与可视化适合作为大数据入门练手或课程设计模板。压缩包共22个文件约4.91MB包含4个Python脚本、2个CSV数据文件、5个HTML可视化结果页、5个XML配置及课程设计报告文档等覆盖从数据读取、清洗到图表输出的完整链路。目前已有3779人学习下载说明其参考价值得到一定验证。读者可获得可直接运行的Spark分析代码、原始与处理后的数据集、逾期率与收入分布等可视化页面以及一份结构完整的课程设计报告便于对照理解项目组织方式、复现分析流程并迁移到自己的数据场景中。1. 信用卡评分卡遇上 Spark一笔 200 万行的样本为什么在单机上跑不动信用卡评分数据分析说白了就是用历史借贷行为去预测「这个人未来会不会逾期」。真正做过的人都知道卡住你的从来不是模型算法而是数据本身——某股份制银行信用卡中心的一份脱敏样本200 万条申请记录、每条约 300 个字段光原始 CSV 就接近 40GB。我最早用 pandas 在 32GB 内存的机器上读read_csv跑到一半直接 OOM改成chunksize分块又发现特征工程里的 groupby 聚合根本没法分块做。这就是 Spark 出场的理由它把数据切成分区并行处理内存放不下就落盘宽表 join 和分组聚合是它的主场。这篇笔记面向两类人一是手里有几十万到上千万条信贷/消费记录、想从零搭一套评分卡流水线的数据工程师二是学过 Spark 语法但没做过完整风控项目的同学。我会按「数据清洗 → 特征工程 → WOE/IV → 逻辑回归评分卡 → 模型评估」这条真实链路走一遍给出能直接抄的 PySpark 代码、参数怎么调、以及我踩过的坑。评分卡本身不复杂难的是让它在 Spark 上跑得又快又对。2. 用 PySpark 把原始信贷流水洗成建模宽表2.1 先想清楚为什么评分卡的数据准备必须上 Spark评分卡建模的数据准备有三个特点决定了它天然适合 Spark。第一是宽表 join申请信息表、还款记录表、征信查询表、额度使用表要按客户号拼在一起pandas 做多表 merge 时中间结果会膨胀好几倍。第二是时间窗口聚合近 6 个月最大逾期天数、近 3 个月查询次数、近 12 个月平均使用率这类滚动特征需要对每个客户按时间切片再聚合单机循环几十万次客户就是灾难。第三是样本量大但特征稀疏真正进模型的特征可能就几十个但候选特征有几百个需要反复试。Spark 的 DataFrame API 和 SQL 都能表达这些操作而且groupBy agg会自动并行。选型上我一般用 PySpark 而不是 Scala风控团队里会 Python 的人多调参和可视化生态也顺性能差距在评分卡这种数据量级上不明显。集群规模上200 万样本、40GB 原始数据我用 4 个 executor、每个 8GB 内存、4 核就能跑得很舒服再大就加 executor 数量而不是单机内存。2.2 环境搭建与读取从 CSV 到带分区的时间序列先确认环境。Spark 3.x 对 Python 3.8 支持最好本地开发用pip install pyspark即可集群上注意 driver 和 executor 的 Python 版本要一致否则会报Python in worker has different version。下面是最小可跑的读取与清洗骨架。from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.types import StructType, StructField, StringType, IntegerType, DoubleType, DateType # 构建 SparkSession评分卡任务内存给足shuffle 分区别太多 spark (SparkSession.builder .appName(credit_scorecard) .config(spark.sql.shuffle.partitions, 200) # 默认 200小集群可降到 50 .config(spark.executor.memory, 8g) .config(spark.sql.adaptive.enabled, true) # 自适应执行Spark 3 必开 .getOrCreate()) # 显式定义 schema避免 inferSchema 全表扫一遍 schema StructType([ StructField(cust_id, StringType(), False), StructField(apply_date, DateType(), True), StructField(age, IntegerType(), True), StructField(credit_limit, DoubleType(), True), StructField(overdue_days, IntegerType(), True), StructField(query_cnt_3m, IntegerType(), True), StructField(label, IntegerType(), True), # 1逾期0正常 ]) df (spark.read .option(header, true) .option(dateFormat, yyyy-MM-dd) .schema(schema) .csv(hdfs:///data/credit/apply_raw.csv)) # 按申请月份分区后续按时间切训练/测试集 df df.withColumn(apply_month, F.date_format(apply_date, yyyy-MM)) df.write.mode(overwrite).partitionBy(apply_month).parquet(hdfs:///data/credit/apply_clean)这段代码有三个关键点。spark.sql.shuffle.partitions默认 200在只有 4 个 executor 的小集群上会产生大量空任务我一般设成 executor 核数的 2~3 倍。spark.sql.adaptive.enabled是 Spark 3 的自适应查询执行能自动合并小分区、优化 join评分卡这种多阶段聚合场景收益明显。显式 schema 比inferSchematrue快很多因为后者要额外扫一遍数据推断类型40GB 数据上能差出好几分钟。2.3 缺失值、异常值与滚动窗口特征三个必做的清洗动作原始信贷数据脏得很具体年龄填 0 或 999、额度为负、逾期天数超过合同期、查询次数缺失。清洗策略要按字段语义定不能一刀切。# 1) 异常值处理年龄限定 18-70额度非负逾期天数不超过 720 df (df .filter((F.col(age) 18) (F.col(age) 70)) .filter(F.col(credit_limit) 0) .withColumn(overdue_days, F.when(F.col(overdue_days) 720, 720) .otherwise(F.col(overdue_days)))) # 2) 缺失值数值型用中位数填充类别型单独处理 median_age df.approxQuantile(age, [0.5], 0.01)[0] # 近似分位数比精确快 df df.fillna({age: median_age, query_cnt_3m: 0}) # 3) 滚动窗口特征近 3 个月查询次数、近 6 个月最大逾期 from pyspark.sql.window import Window w Window.partitionBy(cust_id).orderBy(apply_date).rowsBetween(-5, 0) df df.withColumn(max_overdue_6m, F.max(overdue_days).over(w))approxQuantile用近似算法算中位数误差参数 0.01 表示 1% 精度比percentile_approx更省事大数据量下比精确排序快一个数量级。窗口函数rowsBetween(-5, 0)表示当前行往前 5 行适合按申请次数而非自然月滚动的场景如果要按自然月得先把日期转成月份序号再rangeBetween。这里有个容易翻车的点窗口函数默认要求数据按 partition 和 order 分布如果cust_id数据倾斜严重某个客户几万条记录单个 task 会拖慢整个 job解决办法是对倾斜 key 加盐或改用groupBy 条件聚合。清洗完记得做一次数据质量校验比如标签分布、各字段空值率用df.describe().show()和df.groupBy(label).count().show()快速看一眼别等建模时才发现标签全是 0。3. 特征工程WOE 分箱、IV 筛选与向量化组装3.1 WOE 和 IV 到底在算什么为什么评分卡离不开它们WOEWeight of Evidence和 IVInformation Value是评分卡的灵魂。WOE 衡量某个分箱里「坏人占比」和「好人占比」的差异公式是ln(坏样本占比 / 好样本占比)IV 是各分箱 WOE 的加权和用来衡量一个特征整体的预测力。为什么不用原始值直接进逻辑回归因为信贷特征和违约率往往不是线性关系——年龄 25 岁和 45 岁的违约率可能都低35 岁反而高直接放连续值会让线性模型学不到这种非线性。分箱 WOE 把非线性关系转成单调的数值同时天然处理了缺失值和异常值。IV 的经验阈值小于 0.02 基本没用0.02~0.1 弱预测力0.1~0.3 中等0.3~0.5 强大于 0.5 要警惕数据泄漏比如用了未来信息。我一般保留 IV 大于 0.02 的特征再结合业务常识剔除。3.2 用 Spark 做等频分箱并计算 WOE/IVSpark 没有现成的 WOE 分箱函数得自己写。核心是先算分箱边界再按边界打标最后按箱聚合算 WOE。def calc_woe_iv(df, feature, labellabel, bins5): # 1) 等频分箱用 approxQuantile 取分位点 quantiles df.approxQuantile(feature, [i/bins for i in range(1, bins)], 0.01) # 2) 按分位点切箱边界去重防止空箱 cut_points sorted(set([float(-inf)] quantiles [float(inf)])) expr F.when(F.col(feature) cut_points[1], 0) for i in range(1, len(cut_points) - 1): expr expr.when(F.col(feature) cut_points[i1], i) df df.withColumn(bin, expr.otherwise(len(cut_points) - 2)) # 3) 按箱统计好坏样本 stat (df.groupBy(bin) .agg(F.sum(label).alias(bad), F.count(label).alias(total)) .withColumn(good, F.col(total) - F.col(bad))) total_bad stat.agg(F.sum(bad)).collect()[0][0] total_good stat.agg(F.sum(good)).collect()[0][0] # 4) 计算 WOE 和 IV加 0.5 平滑防止除零 stat (stat .withColumn(bad_rate, (F.col(bad) 0.5) / (total_bad 0.5)) .withColumn(good_rate, (F.col(good) 0.5) / (total_good 0.5)) .withColumn(woe, F.log(F.col(bad_rate) / F.col(good_rate))) .withColumn(iv_part, (F.col(bad_rate) - F.col(good_rate)) * F.col(woe))) iv stat.agg(F.sum(iv_part)).collect()[0][0] return stat, iv # 批量计算所有候选特征的 IV candidates [age, credit_limit, overdue_days, query_cnt_3m, max_overdue_6m] iv_dict {} for feat in candidates: _, iv calc_woe_iv(df, feat) iv_dict[feat] iv print(sorted(iv_dict.items(), keylambda x: -x[1]))approxQuantile的第二个参数是分位点列表第三个是误差。加 0.5 平滑是评分卡的标准做法避免某个箱里坏样本为 0 导致log(0)。这里有个性能坑calc_woe_iv里对每个特征都触发一次approxQuantile和一次groupBy如果候选特征有 200 个就是 400 次 job非常慢。优化办法是把所有特征的分位点一次性算出来或者用QuantileDiscretizerSpark ML 自带批量分箱它内部会合并计算。3.3 把 WOE 值映射回特征并组装成向量算出 WOE 后要把每个样本的原始值替换成对应箱的 WOE 值再用VectorAssembler拼成特征向量喂给逻辑回归。from pyspark.ml.feature import VectorAssembler # 假设已选出 5 个特征逐个映射 WOE selected [age, credit_limit, overdue_days, query_cnt_3m, max_overdue_6m] for feat in selected: stat, _ calc_woe_iv(df, feat) # 构建 bin - woe 的映射广播成字典 woe_map {row[bin]: row[woe] for row in stat.collect()} # 重新打箱并映射 quantiles df.approxQuantile(feat, [0.2, 0.4, 0.6, 0.8], 0.01) cut_points sorted(set([float(-inf)] quantiles [float(inf)])) expr F.when(F.col(feat) cut_points[1], 0) for i in range(1, len(cut_points) - 1): expr expr.when(F.col(feat) cut_points[i1], i) df df.withColumn(feat _woe, F.udf(lambda b: woe_map.get(b, 0.0), double)(expr.otherwise(len(cut_points)-2))) woe_cols [f _woe for f in selected] assembler VectorAssembler(inputColswoe_cols, outputColfeatures) df_vec assembler.transform(df).select(features, label)woe_map用collect()拉到 driver 再广播因为分箱数通常不超过 10字典很小不会撑爆内存。VectorAssembler要求所有输入列是数值型WOE 列已经是 double直接拼即可。注意udf在这里是逐行调用如果数据量特别大可以改用join的方式把stat转成 DataFrame 后按 bin 列 join性能更好但代码稍复杂。我一般数据量在千万级以下就用 udf够用。4. 逻辑回归评分卡训练与评估参数、刻度与验证4.1 逻辑回归在 Spark ML 里的关键参数怎么设评分卡主模型就是带 L2 正则的逻辑回归。Spark ML 的LogisticRegression有几个参数必须调对。from pyspark.ml.classification import LogisticRegression from pyspark.ml.evaluation import BinaryClassificationEvaluator # 按时间切分前 8 个月训练后 2 个月测试避免时间穿越 train df_vec.filter(F.col(apply_month) 2023-08) test df_vec.filter(F.col(apply_month) 2023-08) lr LogisticRegression( featuresColfeatures, labelCollabel, maxIter100, # 迭代上限评分卡一般 50-100 就收敛 regParam0.01, # L2 正则强度越大越防过拟合 elasticNetParam0.0, # 0纯 L2评分卡几乎不用 L1 standardizationTrue, # WOE 已标准化可关掉省时间 threshold0.5 # 分类阈值后续按业务调整 ) model lr.fit(train) # 评估AUC 是评分卡最看重的指标 pred model.transform(test) evaluator BinaryClassificationEvaluator(labelCollabel, metricNameareaUnderROC) print(AUC:, evaluator.evaluate(pred))regParam是防过拟合的核心我一般从 0.01 开始网格搜。maxIter设 100 足够如果 100 次还没收敛多半是特征共线性严重或没做 WOE。standardization在 WOE 特征上可以关掉因为 WOE 本身量纲已经统一。AUC 在信贷场景通常 0.65~0.75 算正常超过 0.8 要怀疑标签泄漏。4.2 把概率转成标准评分PDO、基准分与刻度换算逻辑回归输出的是概率业务要的是 300~850 的整数分。换算公式是score base factor * ln(odds)其中factor PDO / ln(2)base base_score - factor * ln(base_odds)。PDOPoints to Double the Odds表示 odds 翻倍需要的分数常用 20 或 50。import math PDO 20 base_score 600 base_odds 50 # 基准分对应好坏比 50:1 factor PDO / math.log(2) base base_score - factor * math.log(base_odds) # 从模型系数还原每个 WOE 特征的分数 coef model.coefficients.toArray() intercept model.intercept print(factor:, round(factor, 2), base:, round(base, 2)) print(intercept:, round(intercept, 4)) for i, c in enumerate(coef): print(ffeature_{i}: coef{round(c,4)}, score_per_unit{round(-c*factor,2)})score_per_unit表示该特征 WOE 每增加 1 单位总分变化多少。实际部署时把每个特征的「分箱 → WOE → 分数」做成一张映射表线上查表累加即可不需要跑模型。这套换算的好处是分数可解释每降 20 分坏账 odds 翻倍。4.3 用 KS 和 PSI 验证模型上线前必须看的两个指标AUC 只反映排序能力评分卡上线前还要看 KS 和 PSI。KS 衡量好坏样本累计分布的最大差值反映区分度PSI 衡量训练集和测试集或线上新样本的分布差异反映稳定性。# KS 计算按分数排序后累计好坏分布 def calc_ks(pred_df, score_colscore, label_collabel): w Window.orderBy(F.col(score_col).desc()) tmp (pred_df .withColumn(bad_cum, F.sum(F.col(label_col)).over(w) / pred_df.filter(F.col(label_col)1).count()) .withColumn(good_cum, F.sum(1-F.col(label_col)).over(w) / pred_df.filter(F.col(label_col)0).count()) .withColumn(ks, F.abs(F.col(bad_cum) - F.col(good_cum)))) return tmp.agg(F.max(ks)).collect()[0][0] # PSI 计算按分数分箱后比较训练集和测试集占比 def calc_psi(train_df, test_df, score_colscore, bins10): qs train_df.approxQuantile(score_col, [i/bins for i in range(1, bins)], 0.01) cuts sorted(set([float(-inf)] qs [float(inf)])) def bucket(d): expr F.when(F.col(score_col) cuts[1], 0) for i in range(1, len(cuts)-1): expr expr.when(F.col(score_col) cuts[i1], i) return d.withColumn(bucket, expr.otherwise(len(cuts)-2)) t1 bucket(train_df).groupBy(bucket).count().withColumnRenamed(count, c1) t2 bucket(test_df).groupBy(bucket).count().withColumnRenamed(count, c2) n1, n2 train_df.count(), test_df.count() joined t1.join(t2, bucket, outer).fillna(0) psi joined.withColumn(p1, F.col(c1)/n1).withColumn(p2, F.col(c2)/n2) \ .withColumn(psi_part, (F.col(p1)-F.col(p2)) * F.log((F.col(p1)1e-6)/(F.col(p2)1e-6))) \ .agg(F.sum(psi_part)).collect()[0][0] return psiKS 一般要求大于 0.30.2~0.3 勉强可用低于 0.2 模型基本没区分度。PSI 小于 0.1 表示稳定0.1~0.25 需要关注大于 0.25 说明分布漂移严重模型要重新训练。这两个指标我每次上线前必看光看 AUC 会漏掉稳定性问题。5. 避坑与排查评分卡在 Spark 上最容易翻车的 5 个地方5.1 数据倾斜导致某个 task 跑几小时现象job 卡在 99%Spark UI 里某个 task 的 shuffle read 是其他 task 的几十倍。原因cust_id或某个高基数字段做 groupBy/join 时少数 key 的记录数远超平均。解决先df.groupBy(cust_id).count().orderBy(F.desc(count)).show()定位倾斜 key对倾斜 key 加随机盐concat(cust_id, floor(rand()*10))打散后再聚合最后去掉盐或者开启spark.sql.adaptive.skewJoin.enabledtrue让 Spark 自动处理。5.2 WOE 分箱边界在训练集和测试集上不一致现象训练时算好的 WOE 映射表在测试集上映射后 AUC 暴跌。原因分箱边界是用训练集算的测试集如果重新算分位点边界会变导致同一个原始值落到不同箱。解决分箱边界必须只在训练集上计算一次保存成映射表测试集和线上都复用这张表。我一般把cut_points和woe_map一起存成 JSON部署时加载。5.3 用未来信息做特征导致 AUC 虚高现象模型 AUC 0.85上线后实际坏账率完全对不上。原因特征里混入了标签发生之后才知道的信息比如用「当前逾期状态」预测「未来是否逾期」或者滚动窗口包含了未来时间点。解决所有特征的计算时间点必须严格早于标签观察期用Window.orderBy(apply_date).rowsBetween(-N, -1)确保只往前看绝不用rowsBetween(-N, 0)这种包含当前行的写法。5.4 shuffle 分区数不合理拖慢整体速度现象每个 stage 都跑很久Spark UI 里大量小 task。原因spark.sql.shuffle.partitions默认 200小集群上每个分区数据量很小调度开销占比高大集群上又可能分区太少导致单分区过大。解决经验值是 executor 总核数的 2~3 倍。4 executor × 4 核 16 核设 48 左右。开了 AQE 后可以设大一点让它自动合并比如 200。5.5 标签不平衡导致模型偏向多数类现象逾期率只有 2%模型把所有样本都预测为正常准确率 98% 但 AUC 只有 0.5。原因逻辑回归在极端不平衡数据上会偏向多数类。解决不要用准确率评估用 AUC 和 KS训练时对少数类加权Spark ML 的LogisticRegression支持weightCol给坏样本更高权重比如按好坏比 1:10 加权或者对好样本下采样但要注意校准概率。注意评分卡的坑大多不在算法而在数据准备和验证环节。模型跑通只是开始分箱边界、时间穿越、标签泄漏这三件事任何一个出问题上线后都是真金白银的损失。6. 让评分卡跑得更稳增量训练、分数监控与一个我常用的调试习惯模型上线不是终点。信贷客群会漂移经济环境会变一个评分卡通常 6~12 个月就要重新训练。Spark 在这件事上有天然优势历史数据都在 HDFS/数据湖里重新跑一遍流水线就行。但全量重训成本高我一般用增量训练的思路——保留最近 24 个月数据按季度滚动更新每次只替换最老的一个季度。实现上把数据按apply_month分区训练时filter出窗口内的分区即可Spark 的分区裁剪会自动跳过不需要读的文件。分数监控是另一个必须做的事。上线后每周算一次线上样本的分数分布和训练集比 PSI。如果 PSI 连续两周超过 0.1就要拉响警报。下面这个监控脚本我用了很久直接跑在调度平台上。def weekly_monitor(spark, score_table, train_score_path): # 读线上近一周分数 online spark.read.parquet(score_table).filter(dt date_sub(current_date(), 7)) train spark.read.parquet(train_score_path) psi calc_psi(train, online, score_colscore, bins10) ks calc_ks(online, score_colscore, label_collabel) print(fPSI{psi:.4f}, KS{ks:.4f}) # PSI 0.1 或 KS 0.2 触发告警 if psi 0.1 or ks 0.2: print(ALERT: 模型稳定性或区分度下降建议排查) return psi, ks调试评分卡我有个习惯永远先在小样本上跑通全链路再上全量。具体做法是df.sample(0.01)抽 1% 数据把清洗、分箱、训练、评估整条链路跑一遍确认没有语法错误和逻辑漏洞再切全量。这样每次改代码的验证成本从半小时降到两分钟。另一个习惯是把所有中间结果落盘成 parquet比如清洗后的宽表、WOE 映射表、训练好的模型出问题时能快速定位是哪一步的输入变了而不是从头重跑。评分卡这行后悔药就是留好每一步的中间产物。希望帮到你。本文还有配套的精品资源点击获取