简介:这是一份面向大数据与数据分析学习者的完整课程设计资源,主题是基于Spark的信用卡评分数据分析。项目以和鲸社区的信用卡评分模型构建数据为数据集,使用Python调用Spark完成数据预处理、特征分析与可视化,最终形成多张HTML图表,适合想了解Spark实际落地流程或参考课程设计范式的读者。资源共22个文件,zip压缩包约4.91MB,以Python脚本、HTML可视化页面、CSV数据文件和Word课程设计报告为主,其中脚本覆盖数据清洗、分析、Web展示等环节,报告可用于对照整体思路。目前已有3781人学习下载。读者可获得课程设计报告与完整代码,既能借鉴数据导入、预处理和Spark分析的具体写法,也能复用HTML可视化模块,快速搭建自己的信用卡评分数据分析演示项目。
1. 信用卡评分数据分析,为什么要以 Spark 为底座
信用卡评分数据分析,本质上不是在建模,而是在准备能让模型稳定的数据。一个客户过去一年的交易流水、还款行为、额度使用率、逾期天数,这些特征在单机上用 Pandas 也能算,但数据量到千万级甚至亿级交易时,洗数、分箱、做 WOE 转换这些步骤会让单机内存直接见底。Spark 的价值不仅在算得快,更在于它是少数能把“训练时的特征处理”和“上线时的特征处理”用同一套代码跑起来的引擎,这也是它成为信用卡评分场景主流底座的原因。这篇笔记写给做风控数据分析和准备落地评分卡的工程师,按数据准备、特征工程、模型训练、评分转换和坑位排查的顺序,把这条链路讲透。
2. 数据准备与 Spark 环境:第一张表怎么进集群
2.1 为什么评分卡的数据管道绕不开 Spark
信用卡评分和普通的数据分析项目有一个明显区别:它的特征几乎都来自时间窗口聚合。近 6 个月交易笔数、近 3 个月最大逾期天数、额度使用率均值、M1 转 M2 的比例……这些指标都要对每个客户做窗口聚合。单机 Pandas 处理 1000 万条交易流水时,groupby 一次就要几十秒,特征做到几十个维度,迭代一版特征就要等半小时以上。Spark 把数据分散到多个 Executor 上并行聚合,同样规模的流水,分钟级能跑完一轮特征。
另一个常被忽略的原因是训练与上线的一致性。用 Pandas 做特征工程、用 sklearn 训练模型,上线时要么用 Python 重写一套特征逻辑,要么用 PMML 导模型,两边经常对不上数。Spark 的 ML Pipeline 把特征工程和模型训练串成一个 PipelineModel,训练和预测走同一个变换流程,这个一致性对风控模型尤其重要——特征偏移是线上模型翻车的第一大原因。
2.2 最省事的开发环境:本地 Standalone 与 spark-submit 的取舍
很多初学者一上来就搭三节点集群,其实没必要。我做评分卡项目的习惯是:本地开发阶段用 SparkSession 的 local 模式,数据量不大、逻辑能跑通,就立刻切到测试集群验证分布式行为;环境差异主要在资源配置上,代码本身不用改。
from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("credit_card_scoring_dev") \ .master("local[*]") \ .config("spark.sql.shuffle.partitions", "200") \ .config("spark.sql.adaptive.enabled", "true") \ .getOrCreate()这里local[*]表示用本机所有可用核心跑,适合小数据量调试;spark.sql.shuffle.partitions控制了聚合和 Join 后的默认分区数,默认 200,数据量小可以调低到 40~60 减少调度开销,数据量大再往上加。spark.sql.adaptive.enabled是 Spark 3.x 的 AQE,开启后 Spark 会自动合并过小的 shuffle 分区,对倾斜场景有一定缓解作用,建议默认打开。
生产提交时我一般不用pyspark交互式环境,而是把脚本写成文件后用 spark-submit 提交,方便指定队列和资源参数:
spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 8g \ --executor-cores 4 \ --num-executors 20 \ --conf spark.sql.shuffle.partitions=400 \ credit_card_scoring.py--executor-memory和--executor-cores的决定依据是单台物理机的内存和 CPU 配比,常见经验是每个 Executor 内存不超过物理机内存的 1/3,避免堆外内存和系统缓存打架。--num-executors不是越大越好,评分卡特征工程里大量 groupBy 客户号,分区数超出 Executor 数太多会让单个 Executor 串行处理多个分区,GC 压力反而上来。
2.3 schema 先行:信用卡原始数据的读取与类型
信用卡数据分析的原始数据通常来自数仓导出的 CSV 或 Parquet。CSV 读起来方便,但inferSchema会触发一次全文件扫描,而且对日期、金额这类字段经常推断出错。我的习惯是显式定义 schema,既省一次扫描,也避免把时间字段读成字符串、把金额字段读成 double 但丢失精度的尴尬。
from pyspark.sql.types import StructType, StructField, StringType, DoubleType, LongType, DateType customer_schema = StructType([ StructField("customer_id", StringType(), False), StructField("age", DoubleType(), True), StructField("gender", StringType(), True), StructField("occupation", StringType(), True), StructField("annual_income", DoubleType(), True), StructField("credit_limit", DoubleType(), True) ]) df_customer = spark.read \ .option("header", "true") \ .option("sep", ",") \ .schema(customer_schema) \ .csv("hdfs:///data/credit/customer_profile.csv")第三个字段参数False表示非空,Spark 读取时遇到缺失值会直接报错,适合 customer_id 这种主键字段。True的字段允许为空,后续特征工程里统一处理。金额字段用DoubleType时要注意,如果原始数据精确到分,更稳妥的是先用字符串读进来,清洗后再转 DecimalType,避免浮点误差在后续聚合里被放大。
提示:生产环境的数据文件如果有日期分区,建议直接读分区目录,比如
csv("/data/credit/transactions/"),让 Spark 通过分区裁剪只读需要的时间范围,而不是全量读入再过滤。
3. 数据清洗与特征工程:决定评分卡上限的隐藏战场
3.1 信用卡数据的三种脏:缺失、重复与时间泄漏
评分卡项目里 70% 的时间花在特征工程上,特征工程里最容易出问题的不是特征本身,而是数据质量。信用卡数据常见的三种脏:一是缺失,客户职业、收入字段缺失率可能超过 30%;二是重复,同一张卡同一笔交易在数仓里被重复记录;三是时间泄漏,这是最隐蔽的——用未来信息预测过去,模型在验证集上表现极好,上线后立刻崩塌。
我每次拿到数据先跑一个缺失率和重复率检查脚本,把家底摸清楚再动特征。
from pyspark.sql import functions as F df_customer = df_customer.dropDuplicates(["customer_id"]) missing_stats = df_customer.select([ F.mean(F.col(c).isNull().cast("double")).alias(c) for c in df_customer.columns ]).collect()[0] for col_name in df_customer.columns: rate = missing_stats[col_name] if rate > 0.3: print(f"{col_name} 缺失率 {rate:.1%},建议剔除或单独分箱处理")dropDuplicates按客户 ID 去重,这一步必须在聚合之前做,否则同一个客户的记录会被重复聚合,特征值整体偏大。缺失率超过 30% 的字段,常见做法是保留一个“缺失”分箱,而不是直接删除——客户没填收入这件事本身可能就和风险相关。
时间泄漏的排查要结合业务。比如客户当前的逾期标签是 2024 年 6 月定义的,那特征只能用 2024 年 6 月之前的交易流水来算。把交易流水的统计时间窗口截在标签时间之前,是评分卡项目里最基本的纪律。
3.2 WOE 分箱:为什么连续变量不能直接进逻辑回归
信用卡评分卡和普通机器学习模型有一个关键差异:它是给业务人员看的。业务人员要能解释“年龄在 25 到 30 岁之间为什么减分”,这就要求每个特征和分数的关系是单调的、可解释的。直接把年龄、收入这种连续变量丢进逻辑回归,模型本身能跑,但系数解释性差,而且极端值对系数影响很大。
行业标准做法是先分箱,再算每个分箱的 WOE(Weight of Evidence),用 WOE 值代替原始特征进入模型。WOE 公式是:
WOE = ln( (bin_bad / total_bad) / (bin_good / total_good) )它衡量的是这个分箱里坏客户占比相对于整体坏客户占比的偏离程度。WOE 为正说明该分箱坏客户相对集中,为负说明该分箱风险较低。下面的代码演示了按分位数分箱后计算 WOE 的完整过程:
from pyspark.sql import functions as F from pyspark.sql.window import Window as W # 1. 按分位数确定分箱边界 quantiles = df_customer.approxQuantile( "age", [0.2, 0.4, 0.6, 0.8], 0.01 ) bins = [-float("inf")] + quantiles + [float("inf")] # 2. 分箱 df_binned = df_customer.withColumn( "age_bin", F.bucketize(F.col("age"), bins) ) # 3. 按箱聚合好坏客户数 total_bad = df_binned.filter(F.col("default") == 1).count() total_good = df_binned.filter(F.col("default") == 0).count() woe_df = df_binned.groupBy("age_bin").agg( F.sum(F.when(F.col("default") == 1, 1).otherwise(0)).alias("bin_bad"), F.sum(F.when(F.col("default") == 0, 1).otherwise(0)).alias("bin_good") ).withColumn( "woe", F.ln( (F.col("bin_bad") / total_bad) / (F.col("bin_good") / total_good) ) ) woe_df.orderBy("age_bin").show()approxQuantile用近似分位数算法,相比精确排序,在亿级数据上的耗时从分钟级降到秒级,误差通过第三个参数 0.01 控制,即允许 1% 的误差。分箱数量对评分卡影响很大:箱太少丢失信息,箱太多每个箱的样本量不足,WOE 不稳定。常见做法是先用分位数分成 5 箱,再根据业务含义人工合并,比如年龄的边界最后通常会落在 25、30、40、55 这种整数上。分箱完成后,检查每个箱的坏样本数,任何一个箱的坏样本少于 30 个,这个箱的 WOE 就有很强的随机性,需要合并相邻箱。
3.3 从三张表到训练特征:Spark SQL 的特征聚合思路
信用卡评分通常不会只有一张表,而是客户表、交易流水表、还款行为表三张核心表。特征工程的核心是把三张表按客户 ID 聚合到一张宽表上。Spark SQL 在这里比 DataFrame API 直观得多,尤其是涉及多级聚合的时间窗口特征。
CREATE OR REPLACE TEMP VIEW tx_features AS SELECT customer_id, COUNT(*) AS tx_count_6m, AVG(tx_amt) AS avg_amt_6m, MAX(tx_amt) AS max_amt_6m, SUM(CASE WHEN tx_amt > credit_limit * 0.8 THEN 1 ELSE 0 END) AS high_util_cnt_6m, SUM(CASE WHEN tx_time >= DATE_SUB(CURRENT_DATE(), 90) THEN 1 ELSE 0 END) AS tx_count_3m FROM transactions WHERE tx_time >= DATE_SUB(CURRENT_DATE(), 180) GROUP BY customer_id;这里把交易时间窗口拆成 6 个月和 3 个月两档,是为了捕捉近期行为变化。一个客户 6 个月内交易活跃但最近 3 个月突然冷下来,可能是风险信号。DATE_SUB的日期计算在 WHERE 里做,能利用分区裁剪,只扫描近 180 天的数据。high_util_cnt_6m这类特征衡量的是客户是不是经常刷爆额度,在评分卡里通常有较强的区分度。
三张表聚合完成后,以客户宽表为主表做 Left Join 拼接:
df_train = df_customer \ .join(tx_features, "customer_id", "left") \ .join(repay_features, "customer_id", "left")Left Join 会产生大量的空值,比如没有交易记录的客户在交易特征上全是空。这些空值不要直接填 0,而应该单独做一列“是否有交易记录”的二值特征,因为“没交易”和“交易金额为 0”在风险含义上不一样。Join 之后要做一次行数校验,如果结果行数和客户表不一致,说明三张表里有重复的 customer_id,需要回到去重步骤排查。
4. 训练逻辑回归并转换标准评分卡:从系数到分数的最后一步
4.1 用 MLlib 训练带样本权重的逻辑回归
信用卡坏客户占比通常只有 2%~5%,直接把原始样本喂给逻辑回归,模型会把所有客户都预测成好客户,因为全预测“好”的准确率也能到 95%。解决办法是设置样本权重,让模型对坏客户的错分付出更高代价。MLlib 的 LogisticRegression 支持weightCol参数,权重列加在训练数据上即可。
from pyspark.ml.feature import VectorAssembler from pyspark.ml.classification import LogisticRegression from pyspark.ml import Pipeline from pyspark.sql import functions as F # 计算样本权重:让好坏两个类别的权重贡献相等 total = df_train.count() n_bad = df_train.filter(F.col("default") == 1).count() n_good = total - n_bad df_train = df_train.withColumn( "sample_weight", F.when(F.col("default") == 1, total / (2 * n_bad)) .otherwise(total / (2 * n_good)) ) assembler = VectorAssembler( inputCols=feature_cols, outputCol="features" ) lr = LogisticRegression( featuresCol="features", labelCol="default", weightCol="sample_weight", maxIter=100, regParam=0.01, standardization=False ) pipeline = Pipeline(stages=[assembler, lr]) model = pipeline.fit(df_train)standardization=False是评分卡场景里的关键参数。MLlib 默认会对特征做标准化,这会让模型系数变得不可解释——因为系数对应的特征是标准化之后的空间,而上线时我们用的是原始 WOE 值。评分卡的特征经过 WOE 转换后已经在一个相对固定的尺度内,不需要再做标准化,关闭后拿到的系数才能直接用于评分换算。regParam建议从小值开始调,评分卡的逻辑回归不追求极致精度,正则太大系数会被压缩得过平,损失特征区分度。
4.2 评分卡刻度:offset、factor 与一个可抄的公式
模型训练完成后,要把逻辑回归的系数换算成一张评分卡。行业里有一套固定刻度标准:设定基准分 P0、基准好坏比 odds0、翻倍分数 PDO。常见配置是 P0=600 对应 odds0=50:1,PDO=20,意思是好坏比每翻一倍评分增加 20 分。这个配置写进评分卡标准后,全公司不同模型的分数可以直接对比。
换算需要的两个参数 factor 和 offset 来自刻度设定:
factor = PDO / ln(2) offset = P0 - factor * ln(odds0)拿到这两个参数后,把逻辑回归的截距和系数代进下面的公式。这里要特别小心符号方向:MLlib 输出的系数对应的是坏客户概率的 logit,而我们希望分数越高代表风险越低,所以公式里是负号。
import math pdo = 20 # 好坏比翻倍时增加的分数 p0_score = 600 # 基准分 odds0 = 50 # 基准好坏比 factor = pdo / math.log(2) offset = p0_score - factor * math.log(odds0) lr_model = model.stages[-1] intercept = lr_model.intercept coeffs = lr_model.coefficients.toArray() # 对单个样本 x(已是每个特征的 WOE 值),计算总分 def score_from_woe(x: list) -> float: logit_bad = intercept + sum(c * w for c, w in zip(coeffs, x)) return offset - factor * logit_bad这个函数的含义是:先算出坏客户 logit,再用负号把“坏概率高”转换成“分数低”。offset和factor一旦定下来就固定不变,实际生产中可以提前算好两个常量。注意这里的x是 WOE 值而不是原始特征值,如果直接把原始年龄、收入填进来,打出的分是错的。
4.3 验证:用 AUC 与 KS 判断模型可不可用
评分卡模型好不好,不看准确率,看区分度。行业里最常用的两个指标是 AUC 和 KS。AUC 用 MLlib 自带的评估器一行算出来,KS 需要自己按分数分桶计算好坏客户的累计占比差。
from pyspark.ml.evaluation import BinaryClassificationEvaluator evaluator = BinaryClassificationEvaluator( labelCol="default", rawPredictionCol="prediction", metricName="areaUnderROC" ) test_pred = model.transform(df_test) auc = evaluator.evaluate(test_pred) print(f"AUC = {auc:.4f}")AUC 在 0.75 以上说明模型有基本区分能力,0.8 以上在信用卡评分场景里算不错。KS 的计算逻辑是:把样本按预测概率从高到低排序分成 20 桶,算每个桶的累计坏客户占比减去累计好客户占比,取最大值。KS 超过 0.3 模型可用,超过 0.4 区分度很好。注意 KS 和 AUC 都要在测试集上评估,而不是训练集——训练集上的 KS 通常会虚高 0.05 以上。
测试集怎么切也有讲究。随机切分是常见做法,但对时序数据,我更建议按时间切:用前 12 个月的数据训练,后 3 个月的数据验证。原因在下一章展开,随机切分会让验证集和训练集共享同一时间段的信息,模型在验证集上的表现会被高估。
5. Spark 信用卡评分最容易翻车的 5 个细节:现象、原因与解法
5.1 OOM 炸在特征聚合阶段:groupBy 客户号的数据倾斜
现象:跑特征聚合的时候,大部分 Executor 很快跑完,但有一两个 Executor 拖了很久,最后报 OutOfMemoryError,作业整体失败。
原因:信用卡交易数据的客户分布极不均匀,少数大客户可能有几十万条交易流水,普通客户只有几十条。groupBy customer_id 时,热点客户的数据全部落在同一个分区,单个 Executor 扛不住。这不是 Spark 集群资源不够,是数据倾斜。
解决:先看倾斜程度,再决定策略。最常用的做法是加盐(salting):给 customer_id 拼接一个随机后缀,把热点客户的记录打散到多个分区做第一轮聚合,再去掉后缀做第二轮聚合。两阶段聚合的代码会复杂一些,但对评分卡这种高频迭代的任务,值得封装成一个公用函数,一劳永逸。
5.2 小表 Join 也跑得极慢:broadcast 阈值没调
现象:客户表和一张只有几千行的地区代码表做 Join,理论上秒级完成,实际却触发了 Shuffle,跑了十几分钟。
原因:Spark 默认的spark.sql.autoBroadcastJoinThreshold是 10MB,如果小表实际大小超过这个阈值,或者被 AQE 判断为不可广播,就会走 SortMergeJoin,整个过程中所有数据都要重新分区。
解决:在配置里把小表广播阈值调大,比如 100MB,或者在 Join 时显式加 broadcast hint。对评分卡场景来说,地区代码、商户类型这类维度表通常不超过几十 MB,广播出去后每个 Executor 本地保留一份完整副本,Join 完全跳过 Shuffle。
5.3 StringIndexer 在预测时翻车:验证集里多出了新类别
现象:训练跑得很顺利,模型也保存了,但用测试集做 transform 的时候直接抛异常,提示某个字符串列里出现了训练时没见过的类别。
原因:StringIndexer 在训练时会建立一张“字符串 -> 索引”的映射表,预测时遇到映射表里没有的字符串,默认行为是报错。信用卡数据里职业、行业这类字段很容易在验证集里出现训练集没有的新值。
解决:设置StringIndexer的handleInvalid="keep"参数,把未见过的类别统一编码成一个单独的分类。同时要意识到这个新类别本身可能带着风险信号——比如新出现的职业类型,上线时要有对应的监控。
5.4 不设样本权重,逻辑回归把所有客户判成好人
现象:模型训练完,查看测试集预测结果,坏客户的召回率几乎是 0,所有客户都被预测为好客户,AUC 勉强 0.5。
原因:信用卡坏客户占比太低,逻辑回归在没有样本权重的情况下,把全部样本判成好客户就能达到 95% 以上的准确率,模型找不到优化坏客户损失的动力。
解决:按 4.1 节的方法设置weightCol,让坏客户样本的损失权重远大于好客户。设置权重后对比一下好坏两个类别的召回率,如果坏客户召回率上来了但好客户误杀率过高,说明权重过大了,调小比例再跑。这个调参过程没有捷径,是评分卡项目里最依赖经验的环节。
5.5 随机切分验证集,模型上线后效果大跌
现象:离线验证 AUC 0.82,上线后实际区分度只有 0.70,业务部门质疑模型有效性。
原因:随机切分让验证集和训练集来自同一时间段,模型在验证集上相当于“见过类似环境”,而信用卡客群的行为会随时间漂移。经济环境变化、产品政策调整都会让客户行为分布发生变化,随机切分无法暴露这个问题。
解决:改用时间切分。用前 12 个月的数据做训练集,后 3 个月的数据做验证集。这样验证集的样本完全来自模型没见过的未来时间段,更接近真实上线表现。时间切分的另一个好处是可以直接算 PSI(群体稳定性指数),量化训练集和验证集之间的分布漂移程度,PSI 超过 0.1 就要警惕模型上线后的衰减速度。
6. 把评分卡真正落地:两个让模型活下来的进阶习惯
6.1 把 PipelineModel 转成一张分值表,业务系统只需 join
有些风控团队把 PipelineModel 部署成在线服务,但评分卡有一个更贴合行业的落地方式——把模型转换成分值表。评分卡本质上是每个特征每个分箱的分数之和,可以在训练完成后直接生成一张可读的表,业务系统只需要根据客户的各特征分箱去 join 这张表,累加分数即可。这样既不需要在线推理服务,也方便监管审计和业务解释。
score_rows = [] for i, col in enumerate(feature_cols): coef = coeffs[i] # 按特征数平分 offset 和 factor*intercept 的常数部分 base_part = (offset - factor * intercept) / len(feature_cols) for bin_id, woe in woe_map[col].items(): attr_score = base_part - factor * coef * woe score_rows.append((col, bin_id, attr_score)) score_table = spark.createDataFrame(score_rows, ["feature", "bin_id", "score"])这里base_part把常数部分均匀摊到每个特征上,保证所有分箱分数累加后等于标准评分。业务系统拿到客户的原始特征后,先映射到分箱,再 join 这张表汇总,整个打分过程不依赖任何机器学习框架。分值表上线后要注意版本管理,换模型时要同步更新分值表,避免新旧版本混用。
6.2 用 PSI 盯住线上漂移,别让模型悄悄失效
模型上线只是开始。我见过太多评分卡上线时表现不错,半年后区分度下滑到不如规则策略,原因就是客群发生了漂移,但没人发现。PSI 的计算逻辑和 WOE 类似,比较两个时间段各分箱的样本占比差异。
PSI = Σ (actual_ratio_i - expected_ratio_i) * ln(actual_ratio_i / expected_ratio_i)注意:PSI 超过 0.1 说明客群分布发生了明显变化,超过 0.25 说明模型很可能已经失效。
我的习惯是每个月跑一次 PSI 监控,把当月客户的特征分布和建模时的基准分布做对比,同时按月看 KS 和坏账率的实际表现。一旦 PSI 连续两个月超阈值,就要启动模型迭代流程,而不是等到业务投诉才反应。这步没有技术难度,难在坚持——但恰恰是这个习惯,能避免评分卡在无声无息中变成一张废卡。
做这个项目让我最深的教训是:Spark 层面的技术坑都有解,真正让评分卡死掉的是对验证方式偷懒、对上线后漂移忽视。希望这篇笔记能帮你少走几步弯路,祝你的评分卡稳稳落地。
本文还有配套的精品资源,点击获取