简介:这份资源是面向大数据与数据分析初学者、高校课程设计参考者的Spark实战项目,以和鲸社区信用卡评分模型构建数据为数据集,用Python结合Spark完成数据预处理、统计分析与可视化,帮助读者理解分布式框架在真实金融风控场景中的落地流程。压缩包共22个文件,约4.91MB,包含4个py脚本、2个csv数据文件、5个html可视化结果页、5个xml配置及doc课程设计报告等,覆盖代码、数据、图表与文档四类内容,结构完整。目前已有3779人学习下载,说明其参考价值得到一定认可。读者可获得一套可复现的课程设计完整方案,包括数据清洗与特征处理脚本、逾期与收入等维度的可视化页面、项目配置说明以及一份成体系的报告文档,便于对照理解Spark分析流程、复用代码框架或作为同类课题的起步模板。
1. 基于Spark的信用卡评分数据分析:从一份脱敏账单到可解释的评分卡
信用卡评分这件事,真正难的不是模型选型,而是把散落在交易流水、账单周期、还款记录里的行为,压成一张能进风控决策的分数卡。我最早接触这块时,用的是单机 pandas,几十万行还能忍,数据量一上千万行、特征一上几百列,内存直接爆给你看。后来换成 Spark,才把「数据清洗 → 特征工程 → WOE 分箱 → 逻辑回归评分卡」这条链路跑顺。这篇笔记讲的就是这条链路:用 Spark 做信用卡评分数据分析,把原始交易与账单数据变成可解释、可复现、能上线的评分结果。适合两类人——一类是刚接触 Spark、想找一个完整数据分析项目练手的同学;另一类是做风控或商业数据分析、手里有账单类数据但还没跑通分布式流程的从业者。下面按「数据怎么进 → 特征怎么造 → 分数怎么出 → 坑在哪」的顺序讲,每一步都给可抄的命令和参数。
2. 数据接入与清洗:把账单流水变成一张能算的表
信用卡评分的数据源通常有三块:申请信息(人口属性、额度、账龄)、账单汇总(每期账单金额、最低还款、逾期天数)、交易明细(商户类别、金额、时间)。评分卡建模真正吃的是「账户-月份」粒度,也就是每个账户在每个账单周期上的一行特征。所以第一步不是急着建模,而是把这三块对齐到同一个粒度上。
2.1 用 Spark 读多源数据并统一 schema
实际项目里数据格式很杂,CSV、Parquet、JSON 都有。我一般先把原始层落成 Parquet,后面反复读的时候列裁剪和谓词下推都能吃到。下面这段是典型的读入 + 统一字段名 + 类型转换。
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.types import IntegerType, DoubleType, StringType spark = (SparkSession.builder .appName("credit_score_etl") .config("spark.sql.shuffle.partitions", "200") .config("spark.sql.adaptive.enabled", "true") .getOrCreate()) # 账单汇总:账户-账期粒度 billing = (spark.read.parquet("hdfs:///raw/billing/") .select( F.col("acct_id").cast(StringType()), F.to_date("stmt_date", "yyyy-MM-dd").alias("stmt_date"), F.col("bill_amt").cast(DoubleType()), F.col("min_pay").cast(DoubleType()), F.col("past_due_days").cast(IntegerType()), F.col("credit_limit").cast(DoubleType()) )) # 交易明细:需要聚合到账户-账期 txn = (spark.read.parquet("hdfs:///raw/txn/") .select( F.col("acct_id").cast(StringType()), F.to_date("txn_time", "yyyy-MM-dd HH:mm:ss").alias("txn_time"), F.col("txn_amt").cast(DoubleType()), F.col("mcc").cast(StringType()) ))逻辑说明:账单表本身就是账户-账期粒度,直接读;交易表是明细粒度,必须先聚合。spark.sql.shuffle.partitions设 200 是个经验起点,数据量在几千万行级别时够用,太小会导致单分区过大 OOM,太大则小文件过多拖慢 shuffle。spark.sql.adaptive.enabled打开自适应执行,让 Spark 在运行时合并小分区,这个在数据倾斜场景下能救命。
参数说明:stmt_date用to_date显式指定格式,别指望 Spark 自动推断,格式混了会静默变 null。金额字段统一DoubleType,如果对精度要求高(比如对账),换DecimalType(18,2),但计算会慢一些。
2.2 交易聚合与账期对齐
交易明细要按账户和账期聚合,账期用账单日切分。常见做法是把交易时间映射到「所属账单月」,再 group by。
# 把交易时间对齐到账单月:以账单日为界,账单日之后算下期 txn_agg = (txn .withColumn("txn_month", F.date_format("txn_time", "yyyy-MM")) .groupBy("acct_id", "txn_month") .agg( F.count("txn_amt").alias("txn_cnt"), F.sum("txn_amt").alias("txn_amt_sum"), F.avg("txn_amt").alias("txn_amt_avg"), F.countDistinct("mcc").alias("mcc_cnt"), F.max("txn_amt").alias("txn_amt_max") )) # 账单表也补一个月份键,用于 join billing_m = billing.withColumn("txn_month", F.date_format("stmt_date", "yyyy-MM")) panel = (billing_m.join(txn_agg, on=["acct_id", "txn_month"], how="left") .fillna({"txn_cnt": 0, "txn_amt_sum": 0.0, "txn_amt_avg": 0.0, "mcc_cnt": 0, "txn_amt_max": 0.0}))逻辑说明:这里用left join保留所有账单记录,没有交易的月份补 0,而不是丢掉——「这个月没消费」本身就是强特征。countDistinct("mcc")统计消费商户类别数,是衡量消费多样性的常用指标。
参数说明:fillna的默认值要按业务含义给,金额补 0 合理,但如果某个比率类特征缺失,补 0 可能引入偏差,那种情况更适合补中位数或单独加缺失标记列。
提示:join 之前先确认两边的
acct_id有没有前后空格、大小写不一致,这类脏数据在跨系统取数时非常常见,join 不上往往就是它。
2.3 缺失值与异常值处理
清洗阶段最容易被跳过、又最容易翻车的就是异常值。信用卡数据里常见的异常:账单金额为负(退款)、额度为 0、逾期天数超过 999、交易金额出现 6 个数量级的离群点。
panel_clean = (panel .filter(F.col("credit_limit") > 0) .filter((F.col("past_due_days") >= 0) & (F.col("past_due_days") <= 999)) .withColumn("util_rate", F.when(F.col("credit_limit") > 0, F.col("bill_amt") / F.col("credit_limit")).otherwise(None)) .withColumn("util_rate", F.when(F.col("util_rate") > 3, 3.0).otherwise(F.col("util_rate"))) .withColumn("txn_amt_sum", F.when(F.col("txn_amt_sum") > 1e7, 1e7).otherwise(F.col("txn_amt_sum"))))逻辑说明:额度使用率util_rate是评分卡里权重最高的特征之一,先算出来再截断到 3 倍,避免极端值把分箱拉偏。截断阈值不是拍脑袋,一般看分位数,比如 99.9 分位。
参数说明:past_due_days上限 999 是行业里常见的哨兵值约定,超过的当异常处理。截断阈值建议先用approxQuantile看一眼分布再定。
3. 特征工程:用 Spark 造出评分卡真正吃的变量
原始字段直接进模型效果很差,评分卡讲究的是「分箱 + WOE」,把连续变量变成有业务含义的离散段。这一章讲怎么在 Spark 里把特征造出来、分好箱、算好 WOE。
3.1 时间窗口特征与滚动统计
风控里最有区分度的往往是「近 3 个月」「近 6 个月」的行为,而不是当期快照。用 Spark 的窗口函数可以一次算出来。
from pyspark.sql import Window w = Window.partitionBy("acct_id").orderBy("txn_month").rowsBetween(-2, 0) panel_feat = (panel_clean .withColumn("util_rate_3m_avg", F.avg("util_rate").over(w)) .withColumn("past_due_max_3m", F.max("past_due_days").over(w)) .withColumn("txn_amt_3m_sum", F.sum("txn_amt_sum").over(w)) .withColumn("txn_cnt_3m_avg", F.avg("txn_cnt").over(w)))逻辑说明:rowsBetween(-2, 0)表示当前行往前推 2 行,加上当前行共 3 期,正好是近 3 个月。窗口按acct_id分区、txn_month排序,保证时间顺序正确。
参数说明:窗口大小按业务定,3 期和 6 期都常见。注意如果某账户中间有月份缺失,rowsBetween是按行数不是按自然月,严格来说应该先补齐月份序列再算,否则「近 3 期」可能跨了 5 个自然月。
3.2 分箱与 WOE 计算
WOE(Weight of Evidence)是评分卡的核心,衡量每个分箱对好坏样本的区分能力。Spark 里没有现成的 WOE 算子,得自己写。
# 假设 label 列:1 为违约,0 为正常 def calc_woe(df, feature, label="label", bins=10): # 等频分箱 quantiles = df.approxQuantile(feature, [i/bins for i in range(1, bins)], 0.01) cuts = [-float("inf")] + quantiles + [float("inf")] bucket = F.when(F.col(feature) <= cuts[1], 0) for i in range(1, len(cuts)-1): bucket = bucket.when(F.col(feature) <= cuts[i+1], i) bucket = bucket.otherwise(len(cuts)-2) tmp = (df.withColumn("bucket", bucket) .groupBy("bucket") .agg(F.sum(label).alias("bad"), F.count(label).alias("total"))) tmp = tmp.withColumn("good", F.col("total") - F.col("bad")) total_bad = tmp.agg(F.sum("bad")).collect()[0][0] total_good = tmp.agg(F.sum("good")).collect()[0][0] tmp = (tmp.withColumn("bad_rate", F.col("bad") / total_bad) .withColumn("good_rate", F.col("good") / total_good) .withColumn("woe", F.log(F.col("good_rate") / F.col("bad_rate"))) .withColumn("iv", (F.col("good_rate") - F.col("bad_rate")) * F.col("woe"))) return tmp, cuts逻辑说明:先等频分箱,再按箱统计好坏样本数,算 WOE 和 IV。IV 是各箱 IV 之和,用来筛特征,一般 IV 小于 0.02 的特征区分度太弱,可以考虑剔除。
参数说明:approxQuantile的第三个参数 0.01 是允许的相对误差,越小越准但越慢。分箱数 10 是起点,实际会做卡方分箱或决策树分箱来优化。注意 WOE 计算里如果某箱 good 或 bad 为 0,log 会出问题,需要加平滑项。
3.3 特征筛选与相关性检查
造完特征不能全塞进模型,多重共线性会让逻辑回归系数不稳定。用相关系数矩阵筛一遍。
from pyspark.ml.stat import Correlation from pyspark.ml.feature import VectorAssembler feat_cols = ["util_rate_3m_avg", "past_due_max_3m", "txn_amt_3m_sum", "txn_cnt_3m_avg"] assembler = VectorAssembler(inputCols=feat_cols, outputCol="features") vec_df = assembler.transform(panel_feat).select("features") corr = Correlation.corr(vec_df, "features", "pearson").collect()[0][0] print(corr.toArray())逻辑说明:Correlation.corr返回一个矩阵,对角线是 1,非对角线是两两相关系数。一般相关系数绝对值超过 0.7 就考虑去掉一个,保留 IV 更高的那个。
参数说明:VectorAssembler要求输入列都是数值型,类别特征要先做 one-hot 或 WOE 编码。相关系数用 pearson 还是 spearman 看分布,偏态严重用 spearman。
4. 评分卡建模与评估:从逻辑回归到分数映射
特征准备好之后,建模本身反而不复杂,评分卡主流还是逻辑回归,因为可解释。这一章讲怎么在 Spark ML 里训练、评估、把概率转成分数。
4.1 逻辑回归训练与参数设置
from pyspark.ml.classification import LogisticRegression from pyspark.ml.feature import VectorAssembler from pyspark.ml import Pipeline assembler = VectorAssembler(inputCols=feat_cols, outputCol="features") lr = LogisticRegression( featuresCol="features", labelCol="label", maxIter=100, regParam=0.01, elasticNetParam=0.0, standardization=True ) pipeline = Pipeline(stages=[assembler, lr]) train, test = panel_feat.randomSplit([0.7, 0.3], seed=42) model = pipeline.fit(train) pred = model.transform(test)逻辑说明:regParam是 L2 正则系数,防止过拟合,0.01 是个温和起点。elasticNetParam=0表示纯 L2,评分卡一般不用 L1,因为要保留所有特征的系数可解释性。standardization=True对特征标准化,逻辑回归对量纲敏感,这一步别省。
参数说明:maxIter100 通常够收敛,如果没收敛日志会警告,可以加到 200。randomSplit的 seed 固定住,保证每次跑结果一致,方便复现。
4.2 评估指标:AUC、KS 与分数分布
风控不看准确率,看 AUC 和 KS。Spark 自带 AUC,KS 要自己算。
from pyspark.ml.evaluation import BinaryClassificationEvaluator auc = BinaryClassificationEvaluator( labelCol="label", rawPredictionCol="rawPrediction", metricName="areaUnderROC").evaluate(pred) # KS 计算 def calc_ks(pred_df): pdf = (pred_df.select("label", "probability") .rdd.map(lambda r: (float(r[1][1]), float(r[0]))).toDF(["score", "label"])) w = Window.orderBy(F.desc("score")) cum = (pdf.withColumn("cnt", F.count("*").over(w)) .withColumn("bad_cum", F.sum("label").over(w)) .withColumn("good_cum", F.sum(1 - F.col("label")).over(w))) total = pdf.count() total_bad = pdf.agg(F.sum("label")).collect()[0][0] total_good = total - total_bad ks = (cum.withColumn("tpr", F.col("bad_cum") / total_bad) .withColumn("fpr", F.col("good_cum") / total_good) .withColumn("diff", F.abs(F.col("tpr") - F.col("fpr"))) .agg(F.max("diff")).collect()[0][0]) return ks print("AUC:", auc, "KS:", calc_ks(pred))逻辑说明:AUC 衡量排序能力,KS 衡量好坏样本的最大区分度。评分卡项目里 AUC 0.7 以上、KS 0.3 以上算可用,具体阈值看业务容忍度。
参数说明:KS 计算里用了全窗口排序,数据量大时这一步会 shuffle 很重,可以先用approxQuantile分桶再算近似 KS。
4.3 概率转分数:标准评分刻度
业务要的不是概率,是 300 到 850 之间的分数。用标准的 PDO(Points to Double the Odds)公式转换。
import math def prob_to_score(p, base=600, pdo=50, base_odds=50): # base: 基准分, pdo: odds 翻倍所需分数, base_odds: 基准 odds factor = pdo / math.log(2) offset = base - factor * math.log(base_odds) odds = (1 - p) / p return offset + factor * math.log(odds) prob_to_score_udf = F.udf(lambda p: float(prob_to_score(p)), DoubleType()) scored = pred.withColumn("score", prob_to_score_udf(F.col("probability")[1]))逻辑说明:PDO 公式把违约概率映射成整数分数,分数越高信用越好。base=600表示 odds 为base_odds时对应 600 分,pdo=50表示 odds 每翻一倍分数加 50。
参数说明:这三个参数是业务约定,不同机构不一样,建模时要和风控策略对齐,别自己拍。
5. 避坑与排查:评分卡项目里最容易翻车的五件事
这一章是我踩过的坑,按「现象 → 原因 → 解决」写,都是血泪经验。
坑一:join 后数据量暴涨。现象是 join 完行数比左表多好几倍。原因是右表acct_id有重复,或者 join key 有 null 导致笛卡尔积。解决:join 前对右表按 key 去重,或者用left join时先dropDuplicates(["acct_id", "txn_month"]),并检查 key 的 null 比例。
坑二:WOE 计算出现 inf 或 NaN。现象是某箱 WOE 变成无穷大。原因是该箱好样本或坏样本数为 0,log 里出现 0 或除零。解决:加平滑项,比如(good + 0.5) / (total_good + 0.5),或者把样本数过少的箱合并到相邻箱。
坑三:训练集和测试集分数分布差异大。现象是测试集 KS 比训练集低很多。原因是特征里有时间穿越,比如用了未来月份的信息。解决:按时间切分而不是随机切分,训练用早期数据、测试用后期数据,窗口特征严格只用当前及历史月份。
坑四:Spark 任务在 shuffle 阶段 OOM。现象是任务卡在某个 stage 然后 executor 挂掉。原因是数据倾斜,某个acct_id的交易量远超其他。解决:开自适应执行,对倾斜 key 加盐打散,或者把spark.sql.shuffle.partitions调大。用 Spark UI 看每个 task 的 shuffle 读写量,找出倾斜分区。
坑五:分数上线后和离线不一致。现象是离线算的分数和线上实时算的对不上。原因是分箱边界、WOE 映射、缺失值处理在两边实现不一致。解决:把分箱边界和 WOE 表导出成配置文件,线上线下共用同一份,别各写各的。
注意:评分卡项目里,特征口径的一致性比模型精度更重要。一个 AUC 0.75 但口径稳定的模型,比 AUC 0.8 但线上线下对不上的模型有价值得多。
6. 把评分卡跑成可复现的流水线:几个我常用的技巧
到这一步,单次跑通不难,难的是每次换数据、换时间窗口都能稳定复现。我一般会把整条链路包成一个参数化的脚本,用配置文件控制输入输出和分箱参数,而不是改代码。
一个具体技巧是:把分箱边界和 WOE 映射单独落成一张表,建模阶段生成、打分阶段读取。这样线上只需要加载这张表做映射,不用重跑分箱逻辑,既快又不会口径漂移。
# 保存分箱与 WOE 映射 woe_table, cuts = calc_woe(panel_feat, "util_rate_3m_avg") (woe_table.write.mode("overwrite") .parquet("hdfs:///model/woe/util_rate_3m_avg/")) # 打分阶段读取并映射 woe_map = spark.read.parquet("hdfs:///model/woe/util_rate_3m_avg/") # 按 bucket 关联回主表,用 woe 列替换原始值验证方法上,我习惯做两件事:一是用同一份数据跑两遍,确认结果完全一致(排除随机性);二是拿一个已知的坏样本账户,手工走一遍特征计算,和脚本输出对一遍,确认没有逻辑错位。这两步能挡掉大部分低级错误。
参数管理上,把base、pdo、base_odds、分箱数、窗口大小这些全部外置到配置文件,代码里只读不写死。换业务线的时候改配置就行,不用动代码。
最后一个习惯:每次模型迭代都保留一份「模型卡」,记录训练数据时间范围、特征列表、IV 值、AUC、KS、分数分布。过几个月回头看,没有这份记录你根本说不清当时为什么这么定。评分卡这东西,可解释性和可追溯性就是它的命根子,别嫌麻烦。希望帮到你。
本文还有配套的精品资源,点击获取