简介:基于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.x,yarn 模式跑任务,代码用 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/hadoop,docker-compose 直接定义三节点。这里给一个最小可跑的 compose 文件:
version: "3" services: hadoop-namenode: image: bde2020/hadoop-namenode:2.0.0-hadoop2.7.4-java8 environment: - CLUSTER_NAME=spark-cluster ports: - "9870:9870" spark-master: image: bitnami/spark:3.3 environment: - SPARK_MODE=master ports: - "8080:8080" - "7077:7077" spark-worker-1: image: bitnami/spark:3.3 environment: - SPARK_MODE=worker - SPARK_MASTER_URL=spark://spark-master:7077 spark-worker-2: image: bitnami/spark:3.3 environment: - SPARK_MODE=worker - SPARK_MASTER_URL=spark://spark-master:7077这段配置里有个关键点:worker 节点不需要暴露端口到宿主机,它们只要能在内网访问 spark-master 的 7077 端口就行。如果你本机资源紧张,把 worker 数量减到 1 也能跑,只是 shuffle 阶段会退化成单机模式。我一般会额外给 spark-master 和 worker 设置内存上限,bitnami 镜像默认按宿主机可用内存分配,本机 16G 内存跑三节点容易把系统拖垮。可以在 environment 里加 SPARK_WORKER_MEMORY=4g 和 SPARK_DAEMON_MEMORY=1g 限制一下。
容器起来之后,用 spark-submit 提交任务,driver 跑在 master 节点上。这里有一个常见误区:很多人以为 spark-submit 提交后脚本是在本地执行,其实默认 deploy-mode 是 client,driver 跑在提交命令的那台机器上,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", header=True, schema=schema, mode="PERMISSIVE")这里 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 窗口,用 rank=1 去重。
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]), schema="channel 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, boundary="right"): 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_col="is_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 组合。但逻辑上这段代码是对的,在数据几百万行时性能差异不大。
这里有个必须注意的细节:公式里如果某个箱的好样本占比或坏样本占比为 0,log 里会除零得到 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"])] }), schema="cust_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,标准化用 StandardScaler:
from 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(inputCols=feature_cols, outputCol="features_raw") df_vec = assembler.transform(df_train) scaler = StandardScaler(inputCol="features_raw", outputCol="features", withStd=True, withMean=True) 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( featuresCol="features", labelCol="is_default", regParam=0.01, elasticNetParam=0.0, maxIter=100, standardization=False ) 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 映射到指定分数范围。常用基准:设定 odds=1:1 时分数为 600,odds 翻倍时分数增加 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(inputCol="channel", outputCol="channel_idx", handleInvalid="keep") 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 不是一条轻松的路,但一旦跨过临界点,能处理的数据量级和迭代速度会有质的提升。这套链路里最值得投入的其实是特征工程封装和分箱自动化,模型本身反而是最成熟的部分。希望我的这些踩坑记录能帮你把迁移路上的曲折提前绕开。
本文还有配套的精品资源,点击获取