Spark信用卡评分卡分析:从特征工程到WOE与分数映射 简介这是一份基于Spark的信用卡评分数据分析课程设计资源面向大数据或数据分析方向的高校学生与入门开发者适合需要完成类似课程设计或想了解Spark数据处理流程的学习者。资源以和鲸社区信用卡评分模型构建数据为数据集使用Python语言调用Spark完成数据预处理、分析及可视化包含完整的课程设计报告和可运行代码。压缩包共22个文件以Python脚本、HTML可视化图表、CSV数据文件为主另有项目配置文件与PyCharm工程文件整体大小4.91MB。目前已有3781人学习下载。通过该资源可以掌握Spark数据清洗与特征分析的基本思路获得可直接参考的代码框架和可视化结果展示方式对完成课程设计或项目实践有较高参考价值。1. 从一份几千万笔的信用卡流水说起Spark做评分卡分析到底在解决什么如果你正在处理一张几千万笔交易记录的信用卡流水表业务方只给了一个目标——“明天我要看到每个客户的逾期风险评分”那基于Spark的信用卡评分数据分析就是你绕不开的一类Spark数据分析案例。它的核心不是把数据读进pandas里跑个模型而是先把原始流水变成客户级特征再完成WOE分箱、IV筛选、逻辑回归训练和分数映射让最终评分卡能在集群上可复现地生成。这类分析最适合信用卡、消费金融、信贷风控方向的数据分析师和数据工程师。它的价值在两个方面一是数据量变大后单机内存撑不住Spark的分布式聚合能把“按客户汇总几千万笔交易”变成分钟级任务二是评分卡本身要求强解释性Spark SQL配合ML库能让你在同一个Pipeline里完成从数据洞察到模型上线的闭环而不是把数据导来导去。下面这些做法是我在真实风控项目里验证过的照着搭一套最小流程大概需要一台3节点的Spark集群。2. 先把业务翻译成数据信用卡评分分析的特征集与Spark数据预处理信用卡评分分析的第一步不是建模型而是先把“交易流水”聚合成“客户画像”。这里有一个关键认知后续所有分析粒度都是客户不是交易。一笔客户一个月刷了50次卡在评分卡里只应该体现为50次、总金额、平均金额、最大金额、最近一次消费时间等特征而不是50行原始记录。2.1 用Spark SQL把原始流水聚合成客户级特征我会把原始表拆成三份交易流水表、客户基础信息表、逾期标签表。交易流水表保留transaction_id、customer_id、trans_time、trans_amount、merchant_type客户基础信息表包含age、gender、credit_limit、occupation_code等。聚合代码用DataFrame API和Spark SQL都行我习惯先用DataFrame做过滤再用groupBy做聚合。from pyspark.sql import SparkSession from pyspark.sql.functions import count, sum as _sum, avg, max, min, to_date, current_date, date_sub, col spark SparkSession.builder \ .appName(credit_card_score_analysis) \ .config(spark.sql.shuffle.partitions, 200) \ .config(spark.sql.adaptive.coalescePartitions.enabled, true) \ .getOrCreate() tx_df spark.read.parquet(/data/credit_card_tx.parquet) cust_df spark.read.parquet(/data/customer_base.parquet) # 交易时间转成日期方便后续做时间窗口 tx_df tx_df.withColumn(tx_date, to_date(trans_time)) txn_feat tx_df.groupBy(customer_id).agg( count(transaction_id).alias(txn_count), _sum(trans_amount).alias(total_amount), avg(trans_amount).alias(avg_amount), max(trans_amount).alias(max_amount), min(trans_amount).alias(min_amount) )这段代码的关键在于groupBy的key是customer_id最终每个客户只输出一行。千万级交易流水经过groupBy后会被压缩到百万级客户维度这是Spark处理这类分析最舒服的场景——先变大后变小。spark.sql.shuffle.partitions默认就是200如果executor总核心数多可以按实际并发数调整比如32核设成64调太大反而会产生大量小任务调度开销更高。接下来要把近30天交易特征也并进来。注意这里应该先filter再聚合而不是先聚合再过滤。# 近30天交易特征先过滤再做groupBy减少shuffle数据量 recent_30d tx_df.filter(col(tx_date) date_sub(current_date(), 30)) \ .groupBy(customer_id) \ .agg(count(transaction_id).alias(recent_30d_cnt)) # 关联客户基础信息保留全量客户 feat_df cust_df \ .join(txn_feat, customer_id, left) \ .join(recent_30d, customer_id, left) \ .fillna(0, subset[txn_count, total_amount, avg_amount, max_amount, recent_30d_cnt]) # 标签逾期90天以上视为坏客户具体口径要和业务对齐 label_df spark.read.parquet(/data/overdue_label.parquet) feat_df feat_df.join(label_df, customer_id, left) \ .fillna(0, subset[default_flag])join后fillna(0)要特别小心left join产生的null有两种含义一种是“这个客户确实没有交易”另一种是“数据缺失或者join键没对上”。在评分卡里我通常还会加一列has_txn标志位把“没有交易”和“交易量为0”区分开否则模型会把这部分缺失当成负面信号。2.2 做WOE/IV分箱前的字段洞察从describe到直方图特征聚合完之后先别急着建模。评分卡分析里最有价值的环节是看每个特征与好坏客户的关系。我会先用describe()看基本分布再用approxQuantile取分位数作为后续分箱边界最后用Spark SQL算每个分箱的WOE和IV。num_cols [age, credit_limit, txn_count, total_amount, avg_amount] feat_df.select(num_cols).describe().show() # 近似分位数relativeError越小越精确但计算更慢 for c in [total_amount, avg_amount, txn_count]: qs feat_df.approxQuantile(c, [0.2, 0.4, 0.6, 0.8], 0.001) print(c, qs)approxQuantile在Spark 3.x里走的是Greenwald-Khanna近似算法relativeError设0.001已经非常接近精确值。日常探索阶段设0.01就够因为分箱边界初始值后面还要根据WOE趋势手动调整没必要一开始就把精度拉满。拿到分位边界后我习惯先直接用when/otherwise做一版粗糙分箱然后算每个箱子的好客户数和坏客户数再看坏账率是否单调。这一步用Spark SQL最顺from pyspark.sql.functions import when q1, q2, q3, q4 qs # 假设是上面 total_amount 的四分位 feat_df feat_df.withColumn( amount_bin, when(col(total_amount) q1, 1) .when(col(total_amount) q2, 2) .when(col(total_amount) q3, 3) .when(col(total_amount) q4, 4) .otherwise(5) ) feat_df.createOrReplaceTempView(feat_v) spark.sql( SELECT amount_bin, SUM(IF(default_flag1,1,0)) AS bad, SUM(IF(default_flag0,1,0)) AS good, COUNT(*) AS total FROM feat_v GROUP BY amount_bin ORDER BY amount_bin ).show()如果一个特征的分箱结果是箱1坏账率2%箱2坏账率3%箱3坏账率6%箱4坏账率9%那这个特征就有很强的区分力如果坏账率忽高忽低说明分箱边界选得不对需要重分或者把这个特征丢掉。这就是评分卡分析里的“单变量分析”比直接跑模型更容易发现数据质量问题和业务异常。3. 用Spark ML把客户特征变成评分卡从WOE入模到分数映射前面的工作解决了“特征怎么造”这一章解决“怎么把特征变成分数”。评分卡模型最常见的选择是逻辑回归原因不是LR效果最好而是系数可以直接换算成每个分箱的加减分业务和监管都认这个解释性。用GBDT也能出分但解释性成本高很多。3.1 分箱后的WOE替代原始特征为什么评分卡不直接用数值逻辑回归对原始数值的线性假设很强但“总消费金额”和风险之间的关系往往不是线性的——极少消费和高频大额消费都可能是风险信号。所以评分卡的标准做法是把每个连续特征分箱然后用WOE替换原始数值。WOE的定义是某个分箱内坏客户分布占比与好客户分布占比之比的自然对数。计算WOE不需要模型用groupBy加窗口函数就能在Spark上算完from pyspark.sql import Window from pyspark.sql.functions import lit, log grouped feat_df.groupBy(amount_bin) \ .agg( _sum((col(default_flag) 1).cast(int)).alias(bad), _sum((col(default_flag) 0).cast(int)).alias(good) ) w Window.orderBy(lit(1)) woe_df grouped \ .withColumn(total_bad, _sum(bad).over(w)) \ .withColumn(total_good, _sum(good).over(w)) \ .withColumn(bad_dist, col(bad) / col(total_bad)) \ .withColumn(good_dist, col(good) / col(total_good)) \ .withColumn(woe, log(col(bad_dist) / col(good_dist))) \ .withColumn(iv, (col(bad_dist) - col(good_dist)) * log(col(bad_dist) / col(good_dist))) woe_df.orderBy(amount_bin).show()这里Window.orderBy(lit(1))的意思是“全表一个窗口”用它求全局总和避免用groupBy后再join。IV值是WOE加权求和大于0.1的特征可以留下来小于0.02的特征基本没什么区分力。要注意如果某个分箱的good或bad为0WOE会变成inf或null需要做平滑常见做法是每个箱子的bad和good各自加0.5。算完WOE后把WOE表join回feat_df得到每个客户的woe_value特征列。后续建模只用WOE列不再用原始数值。3.2 Spark ML Pipeline特征向量、标准化与逻辑回归训练Spark ML要求输入是向量列所以要先把多个WOE特征列合并成一个VectorAssembler输出列。然后做标准化最后接逻辑回归。from pyspark.ml.feature import VectorAssembler, StandardScaler from pyspark.ml.classification import LogisticRegression from pyspark.ml import Pipeline woe_features [age_woe, credit_limit_woe, txn_count_woe, total_amount_woe, avg_amount_woe] assembler VectorAssembler(inputColswoe_features, outputColfeatures_vec) scaler StandardScaler(inputColfeatures_vec, outputColfeatures_scaled, withStdTrue, withMeanTrue) lr LogisticRegression( featuresColfeatures_scaled, labelColdefault_flag, maxIter100, regParam0.01, standardizationFalse ) pipeline Pipeline(stages[assembler, scaler, lr]) train_df, test_df woe_dataset.randomSplit([0.8, 0.2], seed42) model pipeline.fit(train_df)一个容易被忽略的点StandardScaler的mean和std只能在训练集上fit测试集只能transform。上面这段代码用Pipeline把scaler和lr串在一起fit train_df时标准化的均值和标准差来自训练集transform test_df时用的是训练集参数这就避免了数据泄露。regParam0.01是逻辑回归的L2正则系数如果特征数量多、共线性强可以调到0.05如果AUC没提升也别调太大否则评分卡分箱之间的差异会被压缩。3.3 把概率换算成标准评分Score offset - factor * ln(odds)逻辑回归输出的是坏客户概率。要转成大家习惯的300到700分信用评分需要设定两个业务参数基准分P0和PDO。常见做法是当odds等于1:50时基准分为600分odds每翻一倍分数减少50分。这样因子factor PDO / ln(2)offset P0 - factor * ln(50)。import math from pyspark.sql.types import DoubleType from pyspark.sql.functions import udf def prob_to_score(prob): # 防止概率取到0或1ln(0)会直接报错 p min(max(prob, 1e-6), 1 - 1e-6) odds p / (1 - p) p0 600 pdo 50 factor pdo / math.log(2) offset p0 - factor * math.log(50) return offset - factor * math.log(odds) score_udf udf(prob_to_score, DoubleType()) pred_df model.transform(test_df) scored_df pred_df.withColumn(score, score_udf(col(probability))) scored_df.select(customer_id, probability, score).show(10)这里公式里的负号很关键风险越高、odds越大分越低。如果你发现分数和概率同向涨说明正负号弄反了。分数转换后我一般会再看一下全量客户的分数分布是否近似正态如果明显偏左或偏右说明P0或PDO选得和客群不匹配需要调整基准分。4. 避坑Spark信用卡评分数据分析最常见的5个翻车现场这块内容是血泪经验。评分卡分析看似流程简单但每个环节都有人踩坑。下面5个问题我基本都遇到过按“现象→原因→解决”写给你。4.1 现象join后行数暴涨客户数怎么对不上客户表有500万人交易流水聚合后也有500万个customer_idjoin完却出来700万行。原因几乎都是customer_id有重复或者有一侧为null。空值在join时不会匹配但重复键会让行数翻倍。解决join前先检查键的唯一性。cust_df.groupBy(customer_id).count().filter(count 1).show() txn_feat.groupBy(customer_id).count().filter(count 1).show()如果确实有重复分析业务上保留哪一条比如取最近开户时间那条如果null不是业务合法值直接过滤掉或标记成未知类。还有一个习惯两边都聚合到客户唯一粒度后再join一定不要让明细流水表直接join客户表。4.2 现象WOE计算卡死driver端OOM有人觉得算WOE就是groupBy后collect回本地pandas算几万个分组直接让driver内存爆掉。Spark的根本原则是不要在driver上放数据。正确做法是全程用agg和窗口函数让计算留在executor上。我在3.1节已经写了完整的窗口函数计算法这里再强调Window.orderBy(lit(1))可以拿全局总和不要用groupBy().collect()。如果分箱逻辑非常复杂可以用pandas UDF但returnType必须写清楚否则也会报错。4.3 现象逻辑回归训练结束时loss还在抖AUC很低评分卡样本里坏客户占比通常不到5%逻辑回归很容易收敛到一个“全部预测好客户”的局部最优这时候准确率可能95%以上但AUC只有0.5。另外如果某两个特征高度相关比如total_amount和avg_amountL2正则虽然能压系数但训练过程会变得不稳定。解决一是用AUC而不是accuracy评估二是检查WOE分箱是否存在空箱空箱要合并三是把相关性高的特征做降维或合并。我一般会在训练前用VectorAssembler生成特征向量后用Correlation看相关系数超过0.8的特征只保留IV值更高的那个。4.4 现象Shuffle磁盘溢写严重某个task要跑1小时按customer_id做groupBy时如果某个头部客户有上千万笔交易这个key所在的partiton就会变成数据倾斜点。Spark UI上表现为绝大多数task秒完单独几个task挂在那里。解决思路分两层先看能否在业务上过滤掉异常客户比如识别为商户聚合账户的customer_id如果必须保留可以对customer_id加盐先把每个key随机拆成N个带后缀的子key聚合后再按原始key聚合一次。from pyspark.sql.functions import concat, lit, rand # 第一层聚合加盐后缀拆散热点key salted tx_df.withColumn( salted_id, concat(col(customer_id), lit(_), (rand() * 10).cast(int)) ).groupBy(salted_id).agg( count(transaction_id).alias(part_cnt), _sum(trans_amount).alias(part_amount) ) # 第二层聚合去掉后缀还原成真正客户粒度 result salted.withColumn( customer_id, regexp_extract(col(salted_id), r^(.*)_\d$, 1) ).groupBy(customer_id).agg( _sum(part_cnt).alias(txn_count), _sum(part_amount).alias(total_amount) )加盐的代价是要做两次groupBy所以只在确认有热点key时用如果整体分区数据均匀直接调大spark.sql.shuffle.partitions更划算。4.5 现象pandas UDF返回类型一直报错schema明明看起来对用pandas_udf(returnTypeDoubleType())时函数里如果返回的是pd.DataFrame而不是pd.Series或者返回Series的index没有对齐Spark会在执行阶段报“mismatched schema”。这个问题特别隐蔽因为本地测试时pandas自己不会报错。我的建议是单返回值的场景一律返回pd.Series不要折腾多列StructType确实要多列就在returnType里明确写StructType([...])并且函数内部返回一个单行DataFrame的列表格式。不过说实话在评分卡这个场景里能用agg和窗口函数解决的我不太愿意引入pandas UDF因为一旦变量多了写起来容易翻车调试成本也高。5. 模型验证与业务阈值不是看准确率而是看通过率和坏账率模型训练完不能只看一个指标。评分卡上线前的验证核心是三件事AUC区分度、KS稳定性、阈值通过率。这一章直接讲怎么用Spark把这三件事算出来。5.1 用BinaryClassificationEvaluator看AUC而不是看accuracyclass imbalance严重的时候准确率没有意义。Spark ML自带的BinaryClassificationEvaluator可以直接算AUCfrom pyspark.ml.evaluation import BinaryClassificationEvaluator evaluator BinaryClassificationEvaluator( rawPredictionColrawPrediction, labelColdefault_flag, metricNameareaUnderROC ) auc evaluator.evaluate(model.transform(test_df)) print(AUC, auc)AUC大于0.7说明模型有实用价值0.75到0.8在信用卡评分里已经算不错。如果AUC只有0.55先回去查特征不要调参——多半是某个关键特征没造出来或者标签口径有问题。KS也是风控里常用的单值指标把客户按预测分数从低到高排序计算每个分数档的累计坏客户占比和累计好客户占比之差的最大值。KS越高说明模型能把好坏客户分得更开。pred_df model.transform(test_df) ks_df pred_df.withColumn( score_bin, (col(score) / 50).cast(int) * 50 ).groupBy(score_bin).agg( _sum((col(default_flag) 1).cast(int)).alias(bad), _sum((col(default_flag) 0).cast(int)).alias(good) ) w Window.orderBy(col(score_bin).desc()) ks_df ks_df \ .withColumn(cum_bad, _sum(bad).over(w) / _sum(bad).over(Window.orderBy(lit(1)))) \ .withColumn(cum_good, _sum(good).over(w) / _sum(good).over(Window.orderBy(lit(1)))) \ .withColumn(ks, col(cum_bad) - col(cum_good)) ks_df.select(score_bin, ks).show()这段代码里的score_bin是把连续分数按50分一档离散化目的是让KS曲线更平滑。真正上线时用原始score排序也可以但对几百万客户来说按档位算KS已经够用。5.2 阈值怎么选从分数分布到通过率-坏账率曲线AUC是模型好坏阈值是业务决策。评分卡分析里最后一步通常是确定一个cutoff比如650分以上通过600到650分人工审核600分以下拒绝。Spark上可以用一个聚合查询算出每个分数段下的通过率和坏账率cutoff_df scored_df.withColumn( score_bin, (col(score) / 10).cast(int) * 10 ).groupBy(score_bin).agg( _sum((col(default_flag) 1).cast(int)).alias(bad), count(*).alias(total) ).withColumn(bad_rate, col(bad) / col(total)) cutoff_df.orderBy(col(score_bin).desc()).show(50)如果阈值从700降到650通过率会提升但坏账率也会上升。怎么平衡取决于业务成本和资金计划高风险客群利润率能不能覆盖坏账损失决定了你是卡紧还是放松。这一步不要交给模型自动找最优因为业务代价不是纯统计指标能衡量的。5.3 拒绝推断只用通过客户建模会让样本偏差越来越严重评分卡训练样本来自历史已被批准的客户那些因为评分低被拒绝的客户没有真实表现模型看不到他们的好坏结果。这就导致模型在拒绝客群上的预测能力是“黑匣子”。如果长期只用通过客户迭代模型会陷入一个循环拒得越多样本偏差越大模型对新客户越不准。Spark上做简单拒绝推断有两种常见做法第一种是对被拒绝客户统一打上“未知”标签训练时把它们排除然后给现有通过客户按拒绝比例做加权第二种是给被拒绝客户赋予一个概率权重比如0.5把它们加入训练集让模型至少能看到这部分客群的特征分布。第二种在Spark里只是多了一列sample_weight逻辑回归需要设置weightCol参数。这个方向本身很深不是一篇笔记能讲完但至少你要意识到没有拒绝推断的评分卡上线后通常在低分段表现失真。6. 把整套评分卡分析做成可复用的Spark作业我的几个习惯最后分享三个让我少走弯路的小习惯都来自实际项目管理里的教训。第一个是血缘管理。Spark DataFrame的血缘太长以后任何一个中间步骤出错都要从头开始跑。我会在生成feat_df之后立刻调用checkpoint()把中间结果落盘切断血缘。这样后面做分箱、WOE、建模都是基于稳定数据不会因为改了一行代码就把前面几十分钟的聚合全部重跑。对应的storage level可以设成MEMORY_AND_DISK。第二个是参数固化。评分卡分析最忌讳随机性。randomSplit时固定seedPipeline里所有参数写死在配置文件中每次重跑结果必须一致。我遇到过因为没固定seed上午下午跑出来的AUC差0.02业务方追着问为什么最后发现就是随机数的问题。第三个是保留特征版本。每个分箱的边界、每个WOE值、LR的系数和P0/PDO全部落成一张parquet表存起来。后面模型上线或审计时能回答“这个客户的分数是怎么来的”。这一步其实就是评分卡的“后悔药”没有存档你根本说不清一个分数是哪个版本模型算出来的。这三个习惯让我在接新数据源时能快速复制整套流程而不是每次从零开始调Spark参数。希望帮到你。本文还有配套的精品资源点击获取