简介基于Spark的信用卡评分数据分析课程设计项目面向大数据方向学生、数据分析初学者以及需要课程设计参考的开发者。项目采用Python语言以信用卡评分模型构建数据为数据集完整演示了从数据预处理、特征分析、建模探索到可视化展示的流程适合作为入门Spark批量处理与数据分析的实战案例。压缩包共包含22个文件整体大小约4.91MB内容以Python脚本、HTML可视化页面、CSV数据集和课程设计报告为主其中Python脚本覆盖数据清洗、Web可视化等环节HTML文件可直接查看图表分析结果doc报告则展示完整的方案设计与结论附带项目配置文件以满足复现需要。目前该资源已有3781人学习下载对正在准备大数据课程设计或希望了解Spark数据处理流程的同学具有直接参考价值可节省代码编写与报告撰写时间。1. 从一份 50 万行的信用卡申请数据说起为什么评分卡要换 Spark做信用卡风控的同学对这套流程不会陌生拿申请表、征信报告、历史还款记录拼成宽表跑逻辑回归算 WOE 和 IV最后映射成一张标准评分卡。过去几年我在某公司用 Python Pandas 处理十几万行数据单次跑批还能在十分钟内结束可当数据源切到实时埋点、第三方征信接口和跨月历史明细后单表规模轻松冲到百万到千万行Pandas 的 groupby 和 WOE 分箱开始变得力不从心。一次月度评分卡重建光做特征工程就要跑将近两个小时调一次分箱参数又得从头再来那个阶段我深刻体会到一个事实评分卡流程里最耗时间的不是训练模型而是数据清洗和特征变换。Spark 在这个场景里解决的核心问题不是算法精度而是把「单机内存计算」换成「分布式并行计算」。同样一份 500 万行、200 列的训练宽表用 DataFrame API 做分箱统计、WOE 替换和缺失值填充在四节点集群上能把小时级任务压到十几分钟。本文不讨论 Spark 的基础语法直接从信用卡评分卡的真实落地路径出发讲清楚为什么用 Spark、特征工程怎么写、评分映射怎么做、跑批任务怎么调优以及那些不跑一次根本发现不了的坑。适合的人群是已经在用 Python 做评分卡、想把流程迁移到 Spark 上的风控数据分析师和刚接触 Spark、想用真实业务场景练手的工程师。我会按照一条可复现的链路来讲Spark 环境准备和数据处理 → 分箱与 WOE/IV 计算 → 特征工程落地 → 训练评分卡并映射分数 → 调参与避坑。开始之前先把基础环境立住。我用的方案是本地 Docker 起一个三节点 Spark 集群镜像里自带 Hadoop 和 Spark 3.xyarn 模式跑任务代码用 pyspark 写。你不需要完全一致只要能跑 spark-submit 就行。下面进入正题。2. Spark 环境搭建与信用卡数据加载从 raw 表到训练集的第一个门槛2.1 为什么不用 Pandas 直接做评分卡特征工程先给结论当数据量级在百万行以下、单机内存 32G 以上时Pandas 完全够用但信用卡评分卡场景有个特殊之处——衍生变量极多。一张申请评分卡往往要构造 200~500 个特征包括但不限于近 6 个月平均透支比例、近 3 个月逾期天数最大值、历史贷款申请次数、不同渠道来源的统计聚合。这类特征涉及大量 groupby 和窗口函数Pandas 在千万行表上做 50 个 groupby 操作内存占用会膨胀到原始数据的几十倍因为中间结果反复复制。我曾经在 64G 内存的机器上处理 800 万行数据Pandas 直接 OOM。Spark 的 DataFrame 是惰性求值的所有变换先构建血缘图遇到 action 操作才真正计算。groupby 和窗口函数在分布式环境下走 shuffle 和分片并行内存压力被摊到多台节点上而且 Spark 可以用磁盘溢写兜底不像 Pandas 内存不够直接崩。这不是说 Spark 比 Pandas 快而是说在评分卡这个特征规模大、数据量高的场景里Spark 能把任务跑完。2.2 Docker 起一个三节点 Spark 集群的最小命令本地验证用 Docker 是最快的路径。我日常用的镜像组合是 bitnami/spark 搭配 bitnami/hadoopdocker-compose 直接定义三节点。这里给一个最小可跑的 compose 文件version: 3 services: hadoop-namenode: image: bde2020/hadoop-namenode:2.0.0-hadoop2.7.4-java8 environment: - CLUSTER_NAMEspark-cluster ports: - 9870:9870 spark-master: image: bitnami/spark:3.3 environment: - SPARK_MODEmaster ports: - 8080:8080 - 7077:7077 spark-worker-1: image: bitnami/spark:3.3 environment: - SPARK_MODEworker - SPARK_MASTER_URLspark://spark-master:7077 spark-worker-2: image: bitnami/spark:3.3 environment: - SPARK_MODEworker - SPARK_MASTER_URLspark://spark-master:7077这段配置里有个关键点worker 节点不需要暴露端口到宿主机它们只要能在内网访问 spark-master 的 7077 端口就行。如果你本机资源紧张把 worker 数量减到 1 也能跑只是 shuffle 阶段会退化成单机模式。我一般会额外给 spark-master 和 worker 设置内存上限bitnami 镜像默认按宿主机可用内存分配本机 16G 内存跑三节点容易把系统拖垮。可以在 environment 里加 SPARK_WORKER_MEMORY4g 和 SPARK_DAEMON_MEMORY1g 限制一下。容器起来之后用 spark-submit 提交任务driver 跑在 master 节点上。这里有一个常见误区很多人以为 spark-submit 提交后脚本是在本地执行其实默认 deploy-mode 是 clientdriver 跑在提交命令的那台机器上executor 分散在 worker 节点。本地联调时这样方便看日志但正式跑批建议用 cluster 模式driver 也跑在集群里避免本机断网任务就断掉的尴尬。2.3 信用卡原始表的加载与字段初筛原始数据一般是 CSV 或 Parquet 格式从业务库导出来时常常带着脏数据。我处理过的真实申请数据里身份证号有半角全角混合、手机号有 11 位和带 86 的格式、收入字段有的填 0 有的填空、授信额度有负数。Spark 读 CSV 时如果 schema 推断不准后续类型转换会爆炸所以我不依赖 inferSchema而是手动定义 schema。from pyspark.sql.types import StructType, StructField, StringType, DoubleType, IntegerType, TimestampType schema StructType([ StructField(apply_id, StringType(), True), StructField(cust_id, StringType(), True), StructField(apply_time, TimestampType(), True), StructField(income, DoubleType(), True), StructField(credit_limit, DoubleType(), True), StructField(overdue_days_3m, IntegerType(), True), StructField(loan_cnt_total, IntegerType(), True), StructField(channel, StringType(), True), StructField(is_default, IntegerType(), True) ]) df_raw spark.read.csv(hdfs://namenode:9000/data/credit_apply.csv, headerTrue, schemaschema, modePERMISSIVE)这里 mode 参数选 PERMISSIVE 而不是 FAILFAST因为评分卡场景里单条脏数据直接让整个任务失败得不偿失。PERMISSIVE 模式下解析不了的行会被放进 _corrupt_record 字段后续可以单独检查。如果一行里某个字段类型对不上Spark 会把整行标记为 corrupted 而不是只置空该字段这个行为跟 Pandas 完全不同第一次用很容易踩。更稳妥的做法是读完以后对关键字段做单独清洗和 cast 兜底。数据加载之后第一步不是建模而是看数据的质量报告每列的非空率、唯一值数量、数值分布。Spark 里做这个要小心不要对每一列都调一次 df.describe()那会触发多次全表扫描。我一般把数据 cache 住之后一次性用 agg 函数算多列统计量。from pyspark.sql import functions as F df_raw.cache() agg_exprs [ F.count(F.when(F.col(income).isNull(), 1)).alias(income_null_cnt), F.countDistinct(cust_id).alias(cust_id_distinct), F.min(income).alias(income_min), F.max(income).alias(income_max), F.avg(income).alias(income_avg), F.count(F.when(F.col(overdue_days_3m) 0, 1)).alias(overdue_negative_cnt) ] df_raw.agg(*agg_exprs).show()这段统计在数据量大的时候能明显感觉到 Spark 的优势所有聚合在一次扫描内完成多列并行计算。cache 在这里很关键因为后续还要反复读取数据做特征工程不 cache 的话每次 action 都会重新从 HDFS 读一遍全量数据。cache 的存储级别默认是 MEMORY_ONLY如果数据量大于可用内存部分分片会被重新计算而不是溢写到磁盘这时候建议换成 MEMORY_AND_DISK。数据初筛还有一个容易被忽略的环节删除重复申请。信用卡申请场景里一个客户可能在一个月内提交多次申请同一天重复提交的申请需要按业务规则保留一条。Spark 的 dropDuplicates 可以指定去重键但要注意它只保留第一条出现的记录不保证是业务上最新的那条。我的习惯是先按 apply_time 排序再加 row_number 窗口用 rank1 去重。from pyspark.sql.window import Window w Window.partitionBy(cust_id).orderBy(F.col(apply_time).desc(), F.col(apply_id).desc()) df_dedup df_raw.withColumn(rn, F.row_number().over(w)).filter(rn 1).drop(rn)partitionBy 选 cust_id 是因为一次建模只用每个客户最近的一次申请记录orderBy 里把 apply_time 放前面保证取到最新申请。如果某个客户在一天内多次申请apply_id 的降序排列能兜住时间相同的情况。这个窗口操作在千万级数据上会走 shuffle如果集群资源有限可以先用 groupBy cust_id max(apply_time) 过滤一遍再排序能省掉一半的 shuffle 量。3. WOE 分箱与 IV 计算的 Spark 实现从等频分箱到自动分箱的边界3.1 评分卡里分箱到底在做什么评分卡模型用的是逻辑回归但逻辑回归要求自变量和 logit 之间是线性关系连续型变量往往不满足这个条件。分箱的作用就是把连续变量离散化让每个箱体内的坏样本率接近恒定从而在变量和目标之间建立单调或分段的关系。实际操作中分箱还会顺便处理缺失值和异常值缺失可以单独成箱极端值可以和相邻箱合并避免模型对长尾过度敏感。常见的分箱方法有等距分箱、等频分箱、最优分箱。等距分箱把变量取值范围均分成 N 段实现最简单但容易让大量样本堆在某个箱里等频分箱按分位数切分保证每个箱的样本量基本一致最优分箱则以 IV 值最大或卡方检验的显著性为目标做递归切分。我在 Spark 里实现分箱的时候最常用的是按等频初分、再按坏样本率单调性合并的混合方案先用 approxQuantile 算出 20 个初始切分点然后逐步合并坏样本率趋势不一致的相邻箱。3.2 用 approxQuantile 做等频分箱避开 collect 的坑Pandas 里分箱直接用 qcut 就行但 Spark 的 qcut 并不存在。Spark 提供了 approxQuantile 方法基于 Greenwald-Khanna 算法做近似分位数计算不需要把全量数据收集到 driver这在千万级数据上至关重要。quantiles df_app.groupBy(channel).applyInPandas( lambda pdf: pdf[income].quantile([0.05, 0.25, 0.5, 0.75, 0.95]), schemachannel string, q double )等等这样写不对。applyInPandas 的返回格式不符合预期而且分组后每个分区的分位数计算没有全局意义。正确的做法是全量数据上直接调 approxQuantile如果一定要按渠道分箱那就对每个渠道单独过滤再调用。col_name income quantiles df_app.approxQuantile(col_name, [0.05, 0.25, 0.5, 0.75, 0.95], 0.01) print(approx 5%:, quantiles[0], 25%:, quantiles[1]) # 生成分箱边界 bounds [-float(inf)] quantiles [float(inf)] df_binned df_app.withColumn( income_bin, F.when(F.col(col_name).isNull(), missing) .otherwise(F.array_min(F.transform( F.lit(bounds), lambda b: F.when(F.col(col_name) b, b) ))) )这段代码里有几个关键参数。approxQuantile 的第三个参数 relativeError 设为 0.01 表示允许 1% 的误差值越小精确度越高但计算越慢。分箱后我用了一个复杂的 transform 表达式来给每行打上箱号标签实际生产中这个写法可读性太差我更推荐用 BucketedRandomProjectionLSH 之外更简单的方案——直接用 F.when 链式写法虽然代码长一点但一眼能看懂边界值在哪。实际生产里我一般不用上面那段 transform 写法而是定义一个函数批量生成分箱条件def apply_bins(df, col_name, bin_col_name, cuts, boundaryright): if boundary right: expr F.when(F.col(col_name).isNull(), F.lit(missing)) for i in range(len(cuts)): lower -inf if i 0 else str(cuts[i-1]) upper str(cuts[i]) if i len(cuts)-1 else inf expr expr.when( (F.col(col_name) cuts[i-1]) (F.col(col_name) cuts[i]), F.lit(f[{lower}, {upper}]) ) return df.withColumn(bin_col_name, expr)这段逻辑里最关键的是区间边界的处理方式。用左开右闭区间确保每个值恰好落在一个箱里。缺省值单独成箱不参与区间划分这样后续 WOE 计算时缺失箱可以单独处理。cut 边界值来自 approxQuantile如果某些分位数重复比如 25% 和 50% 分位数相同说明该变量取值稀疏合并箱即可。3.3 WOE 和 IV 的 DataFrame 聚合计算分箱完成后进入评分卡的核心统计环节计算每个箱体的样本总数、坏样本数、好样本数、坏样本率进而得到 WOE 和 IV。WOE 的公式是 ln(坏样本占比 / 好样本占比)IV 是 (坏样本占比 - 好样本占比) * WOE 的加总。Spark 里做这个计算天然适合用 groupBy agg。def compute_woe_iv(df, bin_col, target_colis_default): stats df.groupBy(bin_col).agg( F.count(F.lit(1)).alias(total_cnt), F.sum(F.col(target_col)).alias(bad_cnt), (F.count(F.lit(1)) - F.sum(F.col(target_col))).alias(good_cnt) ) total_bad df.agg(F.sum(F.col(target_col))).collect()[0][0] total_good df.count() - total_bad stats stats.withColumn(bad_pct, F.col(bad_cnt) / total_bad) stats stats.withColumn(good_pct, F.col(good_cnt) / total_good) stats stats.withColumn(woe, F.log(F.col(bad_pct) / F.col(good_pct))) stats stats.withColumn(iv_contrib, (F.col(bad_pct) - F.col(good_pct)) * F.col(woe)) iv_total stats.agg(F.sum(iv_contrib)).collect()[0][0] return stats, iv_total这段代码有一个效率隐患total_bad 和 total_good 两个值各触发了一次全表聚合放在千万级表上等于白白多跑两轮。更优的做法是在同一个聚合里输出全局坏样本数或者先用 cache 过的 df 做一次 count 和 sum 组合。但逻辑上这段代码是对的在数据几百万行时性能差异不大。这里有个必须注意的细节公式里如果某个箱的好样本占比或坏样本占比为 0log 里会除零得到 Infinity 或 NaN后续 Spark 的机器学习库会直接报错或者把模型权重废掉。我的处理方式是给占比做平滑处理加上一个极小值 epsilon 或者直接合并那些占比为 0 的箱体。行业惯例是如果一个箱的坏样本数为 0就把它和相邻箱合并而不是硬算。分箱质量看 IV 值的经验阈值IV 小于 0.02 的变量基本没有预测能力可以剔除0.02 到 0.1 之间是弱变量0.1 到 0.3 是中等强度大于 0.3 要警惕过强变量在评分卡里通常会做额外验证防止变量在未来客群上失效。我一般卡 0.02 的最低线低于这个值的直接不进入建模阶段。3.4 自动分箱的尝试和 Monotonic 约束手动分箱在变量只有几十个的时候还行一旦特征数量超过 100逐个调分箱参数就是灾难。我尝试过在 Spark 里实现决策树式的递归分箱每次选择一个切分点使 IV 增益最大然后用卡方检验判断是否继续分裂。核心逻辑是贪心搜索所有候选切分点对每个候选点计算切分前后的 IV 变化。from pyspark.sql import functions as F def best_split(df, col_name, target_col, candidate_cuts): best_iv -1 best_cut None for cut in candidate_cuts: df_temp df.withColumn( split_flag, F.when(F.col(col_name) cut, 0).otherwise(1) ) _, iv_left compute_woe_iv(df_temp.filter(split_flag 0), split_flag, target_col) _, iv_right compute_woe_iv(df_temp.filter(split_flag 1), split_flag, target_col) split_iv iv_left iv_right if split_iv best_iv: best_iv split_iv best_cut cut return best_cut, best_iv这段代码的效率极差每个候选切分点都触发两次完整的 WOE 计算等于对全表跑几十次聚合。在百万级数据上单变量分箱就要等几分钟100 个变量根本没法用。后来我换了个思路先用 approxQuantile 算出 20 个分位点作为候选切分点然后用 sample 抽取 10% 的数据在内存里用 pandas 做贪心搜索确定最优切分结构最后拿到全量 Spark 上执行分箱。这个方案牺牲了一点精度但把单变量的分箱时间从分钟级压到秒级在实际项目中完全够用。单调性约束在 Spark 里没有现成接口我是靠后处理实现的逻辑回归训练完模型后检查每个分箱对应的系数方向是否一致如果出现某个箱的权重符号和其他箱相反就合并该箱。实际操作中反馈单调性比参数单调性更好验证——用训练好的模型对每个箱打分检查平均分是否随箱号单调变化。4. 特征工程到训练集Spark 管线下从 dirty data 到标准评分卡输入的最后一公里4.1 缺失值填充与异常值截断的 Spark 写法评分卡对缺失值的处理逻辑不是简单填充均值而是分情况完全随机缺失的字段可以填充中位数与目标变量相关的缺失需要用单独的指示变量保留信息。我在 Spark 里的实现方式是先给每一列增加 is_null 指示特征再做填充这样缺失模式本身也进入模型。def fill_missing_with_indicator(df, numeric_cols): df_out df for col in numeric_cols: median_val df.approxQuantile(col, [0.5], 0.01)[0] df_out df_out.withColumn(f{col}_isna, F.col(col).isNull().cast(int)) df_out df_out.withColumn( col, F.when(F.col(col).isNull(), F.lit(median_val)).otherwise(F.col(col)) ) return df_outapproxQuantile 对每个数值列调用一次这会触发 N 次全表扫描。列数多的时候性能极差我通常改成只算一次分位数或者按特征分组减少调用次数。更好的做法是只对缺失率超过 5% 的列做填充缺失率极低的列直接用该列的非空均值填充省去中位数计算。异常值截断在评分卡里用的是 Winsorize 方法把超过 99.5% 分位数的值拉回到 99.5% 分位数。Spark 没有内置 winsorize但可以结合 approxQuantile 和 when 表达式实现lower df.approxQuantile(col_name, [0.005], 0.01)[0] upper df.approxQuantile(col_name, [0.995], 0.01)[0] df df.withColumn( col_name, F.when(F.col(col_name) lower, F.lit(lower)) .when(F.col(col_name) upper, F.lit(upper)) .otherwise(F.col(col_name)) )这里有个经验值信用卡收入字段经常有 0 值0 不是异常但会影响分箱效果。我的处理是先区分真实 0 值和缺失收入为 0 的客户单独成箱不参与 Winsorize。直接用分位数截断会把 0 值也当成有效分布的一部分导致低收入的区分度被压缩。4.2 用 pandas UDF 做复杂特征衍生Spark 的内置函数能覆盖大部分简单变换但有些特征必须用更复杂的逻辑。比如信用卡行为类特征近 6 个月的还款行为中逾期天数从 1 天变成 30 天以上的趋势指标。这类特征涉及跨行比较和模式识别用纯 Spark SQL 写起来非常痛苦。我一般用 pandas UDF 把每组数据拉到一个 pandas DataFrame 里做运算Spark 负责分组和并行。from pyspark.sql.functions import pandas_udf import pandas as pd pandas_udf(double) def calc_overdue_trend(overdue_days_list: pd.Series) - float: if len(overdue_days_list) 3: return 0.0 recent overdue_days_list[-3:].mean() older overdue_days_list[:3].mean() if older 0: return 0.0 return (recent - older) / older df.groupBy(cust_id).applyInPandas( lambda pdf: pd.DataFrame({ cust_id: [pdf[cust_id].iloc[0]], overdue_trend: [calc_overdue_trend(pdf[overdue_days_3m])] }), schemacust_id string, overdue_trend double )pandas UDF 的性能关键在 groupBy 的粒度。按 cust_id 分组时如果每个客户的行数很多每个分组内 pandas 处理的开销可以接受但如果每个组只有两三行UDF 序列化和调用的开销会远大于计算本身。我习惯先确认 groupBy 后的组数占比如果绝大多数组行数小于 5改用 Spark SQL 的窗口函数配合内置聚合效率高一个数量级。4.3 DataFrame 到训练集的转换VectorAssembler 与标准化的坑Spark MLlib 的 LogisticRegression 不接受 DataFrame 多列直接作为特征必须先把所有特征列合并成一个 Vector 列。如果直接喂原始列会报错这一点和 sklearn 完全不同。合并用 VectorAssembler标准化用 StandardScalerfrom pyspark.ml.feature import VectorAssembler, StandardScaler from pyspark.ml.classification import LogisticRegression feature_cols [income_bin_woe, overdue_days_3m_woe, loan_cnt_total_woe, credit_limit_woe] assembler VectorAssembler(inputColsfeature_cols, outputColfeatures_raw) df_vec assembler.transform(df_train) scaler StandardScaler(inputColfeatures_raw, outputColfeatures, withStdTrue, withMeanTrue) scaler_model scaler.fit(df_vec) df_scaled scaler_model.transform(df_vec)这里有两个坑。第一个坑WOE 化之后的特征已经是单变量与目标关系的编码再做标准化其实意义不大逻辑回归对特征尺度不敏感但在正则化时会受影响所以我还是会做标准化保证 L2 正则对每个特征的惩罚是公平的。第二个坑VectorAssembler 不接受 StringType 列作为输入如果直接把分箱标签列放进去会直接抛异常。正确流程是先做 WOE 替换把每个分箱映射成对应的 WOE 值作为数值特征输入模型。WOE 替换的实现是典型的 map 操作把分箱标签映射到 WOE 分数woe_mapping {income_bin: { [0, 5000]: 0.35, (5000, 10000]: 0.12, missing: 0.05 }} def map_woe(df, col_name, mapping): result df for bin_label, woe_val in mapping.items(): result result.withColumn( f{col_name}_woe, F.when(F.col(col_name) bin_label, F.lit(woe_val)) .otherwise(F.col(f{col_name}_woe)) ) return result这段代码看起来笨拙但实际运行效率不差因为 when.otherwise 链不会触发额外扫描所有分支都在同一行内完成。在数据量大时要注意如果映射字典很大生成的表达式会很长Spark 在生成执行计划时可能遇到优化问题。我一般把映射字典限制在 50 个箱以内超出就考虑先 reduce 再 join。4.4 训练验证集的划分与时间窗口陷阱信用卡评分卡不能用随机划分的方式分训练集和测试集因为客群会随时间漂移。正确做法是按申请时间切分用前 6 个月的数据训练最近 1 个月的数据做验证。Spark 的 randomSplit 虽然方便但用在时间序列数据上会产生严重的标签泄漏。我之前一个项目就是随手 randomSplit模型验证集 AUC 高达 0.83上线后实际只有 0.71后来发现验证集里混了大量与训练集高度重叠的客户他们的历史行为已经参与了训练。按时间切分的代码很简单train_df df.filter(F.col(apply_time) F.lit(2024-06-01)) eval_df df.filter(F.col(apply_time) F.lit(2024-06-01))关键在切分日期前先检查目标变量在两个时间窗口内的坏样本率是否发生显著变化。如果坏样本率从 3% 跳到 6%可能是外部环境变化比如政策调整或客群结构变化这时候即使按时间切分模型迁移性也可能很差。我的检查方法是分别统计两个窗口的 mean(is_default)如果绝对差异超过 2 个百分点会考虑缩短训练窗口或者重新定义目标变量窗口。5. 逻辑回归到评分卡映射从概率分数到标准 300-850 分的完整换算5.1 Spark MLlib 逻辑回归的参数选择与训练评分卡建模阶段用的模型并不复杂逻辑回归即可。Spark MLlib 的逻辑回归训练和 sklearn 差异不大但有几个参数需要针对评分卡场景专门调。lr LogisticRegression( featuresColfeatures, labelColis_default, regParam0.01, elasticNetParam0.0, maxIter100, standardizationFalse ) model lr.fit(df_train)regParam 是 L2 正则系数评分卡场景一般取值在 0.001 到 0.1 之间。取值过大会把所有特征系数压向 0变量区分度下降过小则容易过拟合。elasticNetParam 设为 0 表示纯 L2 正则因为评分卡要求保留特征的连续性解释L1 会把部分特征权重置 0不利于业务解释。standardization 这里设为 False因为前面已经做了标准化这里再做一次会重复计算。训练完成后要看系数。逻辑回归的系数和 WOE 编码后的特征配合可以直观看到每个特征对分数的影响方向和大小。这个阶段一个常见的业务校验是逾期天数相关的特征系数应该为正因为 WOE 值越大代表坏样本浓度越高如果系数为负就说明数据里存在辛普森悖论需要回溯分箱逻辑。5.2 从概率到标准分的公式推导评分卡的标准做法不是直接输出模型概率而是通过一个线性变换把 log-odds 映射到指定分数范围。常用基准设定 odds1:1 时分数为 600odds 翻倍时分数增加 20 分。根据这两个锚点可以解出线性变换的参数分数 600 20 * log2(odds / 1)其中 odds p / (1-p)p 为模型预测的违约概率。log-odds 可以由逻辑回归的线性部分直接获得。from pyspark.sql import functions as F df_score model.transform(df_eval) # 包含 probability 列 df_score df_score.withColumn( log_odds, F.log(F.col(probability) / (1 - F.col(probability))) ) df_score df_score.withColumn( score, F.lit(600) F.lit(20 / np.log(2)) * F.col(log_odds) )这里 probability 列实际是向量类型需要用 F.col(probability)[1] 取索引 1 的值作为正样本概率。直接用整列计算会报错这个坑我踩过一次。更稳妥的做法是从模型的 rawPrediction 列取线性部分df_score df_score.withColumn( score, F.lit(600) F.lit(20 / np.log(2)) * (F.col(rawPrediction)[1]) )rawPrediction[1] 和 log-odds 是等价的少一次 log 计算数值稳定性更好。两个锚点的选择会影响分数分布如果希望分数整体抬升可以调整基准分或倍数。行业里常见的评分卡基准分布在 300~850 之间600 分对应 odds 1:1 是一个常见的起点不同公司会按自己的风险偏好调整。5.3 分数校准与业务阈值划定模型算出的原始分是一个相对排序分数不能直接用于业务决策需要做校准。校准方式有两种一种是等比例缩放让历史样本的通过率符合业务预期另一种是重新定义锚点把通过率最高的客群分数拉到一个整数基准。我在项目里用的是第二种因为操作简单且容易向业务解释。# 假设历史数据第 90 分位数的分数是 680 p90_score df_score.approxQuantile(score, [0.9], 0.01)[0] # 希望第 90 分位数对应 700 分 shift 700 - p90_score df_score df_score.withColumn(score_adjusted, F.col(score) F.lit(shift))阈值的划定要看业务成本和收益的平衡审批通过率高但坏账率也高反之则业务增长受阻。我一般会画出不同分数阈值下的通过率和坏账率曲线KS 曲线选择一个「通过率下降速度开始放缓」的点作为审批线。Spark 里做这个分析可以直接对所有分数做累加统计df_ks df_score.groupBy(F.round(F.col(score_adjusted), -2).alias(score_bucket)).agg( F.count(F.lit(1)).alias(cnt), F.sum(is_default).alias(bad_cnt) ).orderBy(F.col(score_bucket).desc()) df_ks df_ks.withColumn( cum_cnt, F.sum(cnt).over(Window.orderBy(F.col(score_bucket).desc())) ) df_ks df_ks.withColumn( cum_bad, F.sum(bad_cnt).over(Window.orderBy(F.col(score_bucket).desc())) ) df_ks df_ks.withColumn(cum_bad_rate, F.col(cum_bad) / F.col(cum_cnt)) df_ks.show(20)这个表出来之后业务就能很直观地看到如果审批线设在 650 分通过率是多少、预期坏账率是多少。把图和数据表交给业务评审比单纯谈模型 AUC 有说服力得多。5.4 评分卡结果导出的 Parquet 与报表格式评分卡上线前需要把每个特征的每个分箱的分数贡献输出成一张明细表供业务审核和监管留档。这张表的字段一般包括特征名、分箱区间、WOE 值、IV 贡献、逻辑回归系数、该箱分数贡献。Spark 里可以用一行加一列的方式生成score_card_items [] for col in feature_cols: for bin_label, woe in woe_mapping[col].items(): coef model.coefficients[feature_cols.index(col)] score_contrib coef * woe * (20 / np.log(2)) score_card_items.append((col, bin_label, woe, coef, score_contrib)) df_score_card spark.createDataFrame(score_card_items, schema[feature, bin, woe, coef, score_contrib]) df_score_card.write.mode(overwrite).parquet(hdfs://namenode:9000/output/score_card.parquet)这里需要注意浮点精度系数和 WOE 相乘之后拿到的分数贡献是一个浮点数导进 Excel 或数据库后四舍五入可能丢掉精度建议保留 6 位小数。另外评分卡明细表必须包含 intercept截距项对应的基准分否则业务无法从明细表独立复算出总分。Intercept 的分数贡献计算方法与普通特征不同是 intercept * (20 / np.log(2))还要加上基准分 600。6. Spark 配置调优与评分卡跑批的避坑指南6.1 内存与并行度的配置经验跑完整个评分卡流程后最影响效率的阶段往往是特征工程里的 groupBy 和窗口操作。Spark 默认的并行度由 total cores 决定但 shuffle 的默认分区数是 200数据量大时容易造成部分分区数据倾斜部分分区数据量极小。我一般会根据数据量显式设置 shuffle 分区数spark.conf.set(spark.sql.shuffle.partitions, 400) spark.conf.set(spark.sql.autoBroadcastJoinThreshold, 10485760)shuffle.partitions 设置为 400 而不是默认 200是因为在四节点集群上每节点 100 个分区能更好地利用 CPU但如果数据量只有几百万行400 个分区反而增加调度开销。经验值是让每个分区处理的数据量在 50MB 到 100MB 之间。autoBroadcastJoinThreshold 的关系不大但如果你有小表要和大表 join调大这个阈值可以让小表广播而不是走 shuffle join在特征映射场景下能省大量时间。数据倾斜是评分卡场景常见的问题。信用卡申请数据里某些渠道来源的样本量可能是其他渠道的 100 倍groupBy channel 后个别 executor 要处理的数据远多于其他 executor整体任务被最慢的 executor 拖死。我的处理手段是加盐把倾斜 key 加上随机前缀打散到多个分区最后再聚合回去。df_skewed df.withColumn( channel_salted, F.when(F.col(channel) hot_channel, F.concat(F.col(channel), F.lit(_), F.floor(F.rand() * 10))) .otherwise(F.col(channel)) ) df_agg df_skewed.groupBy(channel_salted).agg(...) # 去掉盐再聚合一次这段代码里对热点渠道加了 0-9 的随机后缀把一个大 key 拆成 10 个每个分区数据量降到十分之一。加盐的代价是最终聚合要多跑一轮但如果倾斜严重这轮聚合的开销远小于单点故障的等待时间。实际使用中要确认热点 key 的分布如果热点渠道占比超过 30%加盐效果才明显。6.2 cache、checkpoint 和 lineage 的关系Spark 的惰性求值在评分卡流程里有副作用如果某个中间 DataFrame 被后续多个动作重复使用每次动作都会重新计算整个 lineage。比如 df_binned 后面要同时做 WOE 计算、IV 计算和缺失值统计如果不 cache这三次 action 会各自从头执行一遍全部上游变换。数据量大时这是灾难。df_binned.cache() df_binned.count() # 触发真正的计算把数据存到内存cache 之后用 count 触发一次 action 是常用技巧保证后续计算直接读取缓存。如果中间结果太大内存放不下用 checkpoint 把数据写到磁盘并切断 lineage这是一个比 cache 更强的兜底方案。checkpoint 的代价是写盘和读盘的 IO但如果 lineage 特别长比如 50 个步骤checkpoint 反而能加速因为避免了每次从头计算。6.3 三个最常见的跑批异常与定位方法第一个异常是 java.lang.OutOfMemoryError: GC overhead limit exceeded。现象是任务跑到一半 executor 报 OOM重试几次仍然失败。原因一般是每个分区的数据量太大单 executor 内存不足以完成聚合。解决方式不是盲目加内存而是先看 Spark UI 里每个 stage 的 shuffle read 大小如果某个 stage 读了几百 MB 而 executor 只有 2G就把 shuffle.partitions 调大。我遇到最多的情况是窗口函数 partitionBy 的 key 太少数据没被拆散。第二个异常是 org.apache.spark.shuffle.FetchFailedException。现象是某个 executor 在 fetch shuffle 数据时失败通常伴随 executor 失联。原因多半是节点间网络波动或磁盘故障但也有很多次是因为某个分区数据过大导致 executor 在写 shuffle 时把磁盘写满。定位方法是在 Spark UI 看 failed stage 的 input size如果单个文件超过 1GB基本可以确认是数据倾斜。解决方式上文提到的加盐或者把该 key 单独过滤出来走广播。第三个异常不是 Exception而是任务超慢某个 stage 的 duration 是其他 stage 的十几倍。最常见的原因是 Spark 默认把多个 stage 串行化但你在代码里用了多个独立的 action 却没有利用并行。比如先做 A 变量分箱再做 B 变量分箱如果这两个操作没有任何依赖关系可以用异步提交的方式并行执行。我通常会把流程拆成多个子任务分别提交而不是写成一个巨大的脚本从头跑到尾。6.4 一个隐蔽的坑当评分卡特征里有 text 列时 VectorAssembler 的行为信用卡申请数据里常常有职业、行业、居住城市这类文本字段。这些字段不能直接作为特征进入逻辑回归必须要做编码。Spark 的 StringIndexer 会按照频率给类别编码但这个方法对低频类别不友好新数据里出现训练集未见的类别时会有 handleInvalid 参数控制行为默认是 error 直接报错。from pyspark.ml.feature import StringIndexer indexer StringIndexer(inputColchannel, outputColchannel_idx, handleInvalidkeep) df_indexed indexer.fit(df_train).transform(df_train)handleInvalid 设成 keep 会把未见过的类别编码到一个独立索引而不是报错或置为 null。这在上线阶段非常重要因为真实业务流量里随时可能出现新渠道。但要注意keep 模式下未来新类别的编码值是一个未知数可能落在任何位置模型解释性会受到一点影响。更稳妥的做法是把类别频率低于某个阈值的值统一合并成 other 类别再走索引。我见过一个项目因为某个城市名拼写错误导致新类别出现全部样本跑出 NaN 预测概率后来排查了整整一个下午才发现是 StringIndexer 的 handleInvalid 默认行为在搞鬼。这类问题在评分卡上线运维中非常常见代码逻辑本身没有错但数据分布变化让模型失效了。7. 最后一章评分卡迁移到 Spark 之后我用哪些方法验证结果没跑偏整个流程跑通之后最需要回答的问题是Spark 算出来的评分卡结果跟原来 Python 单机算出来的比有没有系统性偏差这个验证不能靠感觉要有量化对比。我的做法是准备一份小样本数据集大约十万行分别用 Pandas 流程和 Spark 流程做完整评分卡建模对比三个关键指标各特征 IV 值相关性、逻辑回归系数相关性、最终评分的分位数分布。评分分位数分布很容易算# 对比 Pandas 和 Spark 两条流程产出的分数分位数 spark_quantiles df_score.approxQuantile(score, [0.1, 0.25, 0.5, 0.75, 0.9], 0.01)理论上如果特征工程逻辑一致、分箱边界一致、WOE 映射一致两条流程的分数分位数差异应该小于 5 分。如果差异超过 10 分优先检查分箱边界是否一致——尤其注意 approxQuantile 的近似误差在不同数据量级下表现不一样小样本上误差更大。验证通过之后一个重要的工程习惯是给整个评分卡管线的每一步生成一个摘要输出并落盘。比如分箱明细表、WOE 表、模型系数表、分数分布表每个步骤的结果单独存成一个 parquet 文件日期作为分区字段。这样一旦业务反馈某天的分数分布异常可以直接回溯到对应日期的分箱表和数据质量报告定位是数据源变了还是分箱参数被意外改动。我的习惯是每天跑批后自动对比当天和前一天的分位数分布差异超过阈值直接报警。从 Pandas 迁移到 Spark 不是一条轻松的路但一旦跨过临界点能处理的数据量级和迭代速度会有质的提升。这套链路里最值得投入的其实是特征工程封装和分箱自动化模型本身反而是最成熟的部分。希望我的这些踩坑记录能帮你把迁移路上的曲折提前绕开。本文还有配套的精品资源点击获取