简介:一套基于Hadoop与Spark的大数据金融信贷风控系统设计与实现源码包,面向计算机相关专业毕设、课设及大数据技术进阶者。系统围绕信贷风控场景,覆盖数据采集、流式处理、风险识别等环节,内含Java与Scala编写的核心业务逻辑和Spark Streaming数据源模块,并结合XML配置、Properties配置项及SQL建表脚本,便于理解项目结构并快速搭建运行环境。压缩包共69个文件,整体仅58KB,以源码为主、结构紧凑。其中36个Java文件对应业务处理与接口实现,8个Scala文件负责Spark计算逻辑,12个XML文件用于框架配置,另有SQL脚本与Markdown文档辅助数据库初始化和项目说明。已有325人浏览学习,代码经测试运行成功,可直接用于毕设演示、课程设计或二次开发。资源保留了完整工程目录(如credit-risk-control、data-source-spark-streaming等),并含README说明,为学习大数据生态组件在真实风控项目中的整合提供了有价值的参考。
1. 这个系统到底在解决什么问题
先说一个反直觉的结论:这类信贷风控项目翻车,多数不是算法不行,而是数据管道先垮了。一个基于Hadoop、Spark实现的大数据金融信贷风险控系统,本质上干的是三件事——用HDFS把千万级申请记录、还款流水、用户画像存下来;用Spark把特征工程和批量评分跑起来;再用Spark MLlib训练出能解释给业务听的风险模型。很多课程设计和入门项目把精力都花在调模型上,结果一到千万行样本、上百列特征,单机Python直接内存溢出,Hive跑一轮特征要几个小时,这时才回头补存储和计算的功课。这套系统适合三类人:做大数据方向课程设计的学生、中小信贷或助贷平台想自建风控引擎的数据开发、以及从数据工程往算法方向转的从业者。它解决的核心问题不是“模型更聪明”,而是“在数据量上来之后,风控链路还能按天稳定跑完”。
2. 整体架构与数据管道:先想清楚数据怎么流,再动代码
2.1 三层架构与选型理由
常见的落地做法是把系统拆成三层:数据接入层、存储计算层、模型应用层。接入层接收借款申请、还款流水、用户授权的外部征信数据,通过Kafka或者离线文件批量落盘;存储计算层以HDFS为数仓底座,用Spark跑ETL和特征加工;模型应用层把训练好的模型和评分结果提供给信贷审批接口调用,同时会落到运营看板(也就是常说的数据大屏)做风险监控。
为什么不用业务数据库直接算?信贷风控的特征工程往往要回溯6到12个月的历史数据,业务库通常只保留最近几个月,而且十几张表关联聚合的查询会把在线库拖垮。HDFS天然适合全量快照和长周期存储,做了数据分层之后,每一层都在上一层的成本上增量计算,调试时也能按分区回溯。
计算层选Spark而不是纯Hive MapReduce,原因是特征加工涉及大量窗口聚合和多表关联,MapReduce每跑一个Stage都把中间结果落盘,一轮全量特征跑下来小时级起步;Spark把中间结果留在内存,同样逻辑能压缩到十几分钟。代价是内存参数调不好就会踩到一堆Spark内存相关的坑,这部分在第五章单独讲。至于实时流计算,信贷审批里的反欺诈确实有实时场景,但这个标题的主线是批处理,T+1跑批已经覆盖大部分信贷风险管理诉求,实时模块后面按需再接Flink更稳妥。
2.2 HDFS目录规划与分区策略
目录规划不要等数据进来了再改,起步就按分层建好。我一般会这样组织:
hdfs dfs -mkdir -p /data/credit/ods/{loan_apply,repay_plan,user_profile} hdfs dfs -mkdir -p /data/credit/cds/{feature,label} hdfs dfs -mkdir -p /data/credit/ads/{score,risk_report}ODS层存原始数据,CDS层是清洗后的宽表和标签,ADS层对外输出风险分和报表。表内统一按业务日期做分区,比如dt=2024-06-30,存储格式用Parquet加Snappy压缩。这样规划的好处是:按天分区让任务天然支持增量跑批,补数时只重算对应分区;Parquet列式存储配合压缩能把单日几十GB的交易明细压缩到四分之一左右,读特征时只需扫描相关列,对小集群的磁盘IO压力小很多。
| 分层 | 表名示例 | 内容定位 | 保留周期 |
|---|---|---|---|
| ODS | ods.loan_apply | 原始申请流水,按天分区 | 全量保留 |
| CDS | cds.loan_feature | 清洗去重后的客户特征宽表 | 全量保留 |
| ADS | ads.credit_score | 模型打分结果与风险标签 | 最近24个月 |
还有一点金融项目里要养成习惯:身份证号、手机号这类敏感字段在写入ODS时就做MD5或AES脱敏,模型训练用脱敏后的ID做关联即可,尽量避免明文敏感数据长期躺在HDFS上。这个意识在课程设计和生产环境同样重要,别等安全的同学来找你。
2.3 Spark ETL的核心参数:并行度、内存与Shuffle
ETL跑不跑得动,一半取决于SparkSession的参数。这是我在风控项目里最常用的一段启动配置:
from pyspark.sql import SparkSession spark = (SparkSession.builder .appName("credit_risk_etl") .enableHiveSupport() .config("spark.sql.shuffle.partitions", "200") .config("spark.sql.autoBroadcastJoinThreshold", "10485760") .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") .getOrCreate())spark.sql.shuffle.partitions是最容易忽略的参数。默认200在数据量小时没问题,但数据量上来后,每个分区处理的数据不均会产生大量小任务。经验值是设成集群总vcore数的2到3倍,比如20个executor、每executor 4核,总vcore 80,设200左右比较合理。autoBroadcastJoinThreshold设为10MB,小于这个尺寸的维表自动走Broadcast Join,避免触发Shuffle。KryoSerializer能显著减少序列化开销,但要注意模型类需要注册Kryo类,否则退化到Java序列化反而更慢。
特征加工的核心逻辑是滑动窗口聚合。下面这段SQL是典型的客户行为特征:
INSERT OVERWRITE TABLE cds.loan_feature PARTITION(dt='2024-06-30') SELECT user_id, COUNT(*) AS apply_cnt_180d, SUM(loan_amt) AS total_loan_amt_180d, AVG(credit_score) AS avg_credit_score, SUM(CASE WHEN overdue_days > 30 THEN 1 ELSE 0 END) AS overdue_cnt_180d FROM ods.loan_apply WHERE dt BETWEEN date_add('2024-06-30', -180) AND '2024-06-30' GROUP BY user_id窗口取180天而不是30天的原因,是信贷风险的特征衰减周期远长于支付或营销场景,一个客户半年前的一次逾期对当前还款意愿仍有明显影响。spark.read.json在接入外部数据时很常用,原始JSON文件先读进来再写成Parquet,后续查询性能能差一个数量级。
3. 风险模型训练:选型、样本不均衡与参数调优
3.1 为什么用Spark MLlib而不是单机Python
很多从数据分析转过来的人习惯Pandas加Scikit-learn,这套组合在特征宽表达到几百万行、几百列时就会出问题。Pandas的DataFrame在复制和过滤时会放大内存占用,特征宽表几十GB,单机内存直接打满;Sklearn的逻辑回归训练要反复扫描全量数据,千万级样本一次迭代就是几分钟,做网格搜索基本不可行。
MLlib的核心优势是数据并行。训练数据分片存放在executor上,每次迭代聚合梯度时只传梯度向量而不是全量数据,千万级样本的LR训练能被压到分钟级。当然MLlib的算法实现比Sklearn朴素一些,没有自动早停和丰富的求解器选择,但风控场景的诉求是“能按天稳定跑批”,分布式训练的价值远大于单机上的算法花活。
这也是你在简历或答辩里会被追问的地方:为什么不用XGBoost?答案是风控模型要能解释,XGBoost可以做辅助,但主线模型是评分卡,这个选型逻辑在下一节展开。
3.2 评分卡还是随机森林:可解释性优先
信贷风控里模型的可解释性不是加分项,是硬要求。拒绝一笔借款要能给客户和审核人员讲清楚“为什么拒绝”,监管检查时也要能说明每个变量的影响方向。这两种模型的取舍可以看这张表:
| 维度 | 逻辑回归评分卡 | 随机森林/GBT |
|---|---|---|
| 可解释性 | 每个特征有明确权重,可转标准评分 | 黑匣子,只能给特征重要性排序 |
| 非线性关系 | 需人工做分箱和交叉特征 | 自动捕捉非线性 |
| 训练成本 | 低,分钟级 | 较高,树模型调参空间大 |
| 上线接受度 | 业务和监管都认可 | 客户申诉时难解释 |
| 典型用途 | 主评分卡 | 变量筛选、模型上限探测 |
常见的做法是双轨并行。先用逻辑回归评分卡做主线版本,保证可解释和可上线;再跑一版随机森林或GBT,看同样的特征下区分度上限有多高——如果树模型AUC比LR高出一大截,说明特征工程还没做透,回去补特征而不是换模型。随机森林的特征重要性还可以反过来指导LR的特征筛选,把重要性低的变量剔除,LR的稳定性通常会更好。
进阶一些的团队会把两者结合,用GBDT的叶子节点做特征编码喂给LR,也就是业界说的LR加GBDT组合模型。这个方案保留了LR的输出层,可解释性打了折扣,但区分度确实能上一个台阶。对于课程设计或第一版系统,先把评分卡跑通,组合模型作为二期优化方向。
3.3 训练代码落地与参数设置
PySpark的风控训练代码结构很固定,下面这段是核心流程:
from pyspark.ml.feature import VectorAssembler, StandardScaler from pyspark.ml.classification import LogisticRegression from pyspark.ml.evaluation import BinaryClassificationEvaluator feature_cols = ["age", "income", "loan_amt", "credit_score", "apply_cnt_180d", "overdue_cnt_180d"] assembler = VectorAssembler(inputCols=feature_cols, outputCol="features_vec") scaler = StandardScaler(inputCol="features_vec", outputCol="features") lr = LogisticRegression( featuresCol="features", labelCol="label", maxIter=100, regParam=0.01, elasticNetParam=0.0, standardization=False )VectorAssembler把多列特征拼成一个向量,StandardScaler做标准化,让逻辑回归的梯度下降收敛更快。这里要特别注意standardization=False:数据已经手动标准化过,再开LR自带的标准化等于做了两遍,系数解释会乱掉。elasticNetParam=0.0表示纯L2正则,风控场景一般不用L1,因为L1会把特征系数压成0,不利于评分卡解释。regParam取0.01是起点,如果训练集AUC和验证集AUC差距超过0.05,说明过拟合,把正则系数往上调。
标签的定义是整个模型的地基。信贷里常用M3逾期作为坏客户定义,也就是逾期超过90天。构造标签时要注意观察期和表现期严格错开:
-- 用T时刻的特征,预测T+90天内是否发生M3逾期 SELECT user_id, CASE WHEN max(overdue_days) >= 91 THEN 1 ELSE 0 END AS label FROM ods.repay_plan WHERE due_date BETWEEN date_add('2024-06-30', 1) AND date_add('2024-06-30', 90) GROUP BY user_id这段SQL的含义是:在2024年6月30日这个观察点上,看未来90天内这个客户是否逾期超过90天。如果特征里已经包含了未来信息,比如把T+90天的还款行为算进特征,那就是标签泄漏,离线AUC会虚高到0.9以上,上线后完全不收敛,这是风控建模里最隐蔽的翻车点。
3.4 样本不均衡的处理策略
信贷场景的坏样本占比通常只有2%到5%,直接训练的话模型会倾向于把所有样本都预测为“好客户”。常见的处理方式有三种:负样本过采样、正样本下采样、损失函数加权。在分布式环境里,SMOTE这类插值方法实现起来比较麻烦,最常用的是对坏样本做重复采样,或者给坏样本加大权重。
# 下采样:把好样本抽到坏样本的5倍左右 from pyspark.sql import functions as F df_good = df.filter(F.col("label") == 0).sample(withReplacement=False, fraction=0.3, seed=42) df_bad = df.filter(F.col("label") == 1) df_balanced = df_good.union(df_bad)采样比例取多少需要反复试,5比1是一个比较稳妥的起点。但要注意,采样会扭曲原始分布,模型输出的概率不代表真实违约概率。上线时不能直接用0.5作为阈值,而是要在原始分布上重新校准阈值——后面第六章会讲这个校准方法。另一个思路是给模型传classWeight,PySpark的逻辑回归支持设置weightCol,给坏样本赋5到10倍权重,这样不用改数据分布,阈值校准也更简单。课程设计里两种都做一遍,对比一下区别,答辩素材就有了。
4. 源码组织、依赖版本与运行方式:从伪分布式到集群
4.1 源码结构与文档说明
标题里既然带了“源代码+文档说明”,这部分就得按一个能交付的标准来组织。目录结构大致是这样:
risk-credit-engine/ ├── conf/ │ ├── application.yaml │ └── hdfs-site.xml ├── docs/ │ ├── 设计文档.md │ └── 部署文档.md ├── etl/ │ ├── ods_to_cds.py │ └── feature_engineer.py ├── model/ │ ├── train_scorecard.py │ └── evaluate.py ├── deploy/ │ ├── spark_submit.sh │ └── schedule.sh └── pom.xml我的习惯是强制要求配置不硬编码。HDFS路径、数据库连接、模型参数全部放conf/application.yaml,脚本启动时读取,换环境只改配置不改代码。docs/设计文档.md写清楚数据从哪来、特征怎么定义、标签怎么生成、模型怎么评估;部署文档.md写清单——需要什么版本的JDK、Hadoop、Spark,哪些Jar包要放到哪个节点。别小看这两个文档,半年后你自己回来看项目,靠的就是它们。
4.2 依赖声明与版本配对
如果工程用Java或Scala写,pom文件的依赖管理是第一个坑。下面这段是Spark 2.4.x时期的经典依赖配置:
<dependencies> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-core_2.11</artifactId> <version>2.4.8</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-sql_2.11</artifactId> <version>2.4.8</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.spark</groupId> <artifactId>spark-mllib_2.11</artifactId> <version>2.4.8</version> <scope>provided</scope> </dependency> </dependencies>三个关键点:scope必须设成provided,因为集群上已经有Spark环境,打包进去反而会和集群内的Hive、Parquet版本冲突;Spark 2.4.x对应Scala 2.11,Spark 3.x对应Scala 2.12,依赖名里的_2.11后缀要跟着改;Hadoop client的版本必须和集群一致,否则运行时会出现NoSuchMethodError。用官方预编译版的Hadoop和Spark,本机设置好JAVA_HOME和HADOOP_HOME就能跑起来,不用自己编译源码,这是最省心的路径。
4.3 从伪分布式调试到集群提交
很多初学者一上来就搭三节点集群,折腾两天环境没起来,代码一行没跑。我建议先单机用Hadoop伪分布式模式把链路打通。伪分布式就是在一台机器上同时启动NameNode、DataNode、ResourceManager、NodeManager四个进程,数据量小、排错快,代码逻辑问题在本地暴露得最直接。核心是改三个配置:core-site.xml指定NameNode地址,hdfs-site.xml设副本数为1,yarn-site.xml配置ResourceManager地址,启动后jps能看到四个进程就说明环境起来了。
集群环境跑批任务的提交命令大概是这样的:
spark-submit \ --master yarn \ --deploy-mode cluster \ --queue risk_engine \ --num-executors 20 \ --executor-cores 4 \ --executor-memory 8g \ --driver-memory 4g \ --class com.risk.batch.TrainMain \ risk-batch-1.0.jar \ --run-date 2024-06-30--queue指定Yarn队列,把风控任务和普通数仓任务隔离,避免互相抢资源。--num-executors乘--executor-cores得到总并发度,要控制在Yarn队列上限以内。--executor-memory不是越大越好,超过16g容易触发GC长暂停,后面避坑章节会细说。本地调试时把--master换成local[*],小数据量直接跑通流程再上集群。
5. 信贷风控系统落地的5个避坑记录
5.1 标签时间窗重叠导致AUC虚高
现象:离线验证集AUC达到0.92,团队信心满满上线,三个月后复盘发现坏账率没有任何下降。原因:构造标签时把观察期和表现期重叠了。比如用T时刻的申请数据做特征,却把T+30天的逾期结果也算进了特征里,模型在离线测试时“看到”了未来信息。解决:特征和标签严格按时间切分,特征只使用T日之前的数据,标签只看T+1到T+90的表现期。这个检查要写进数据质量校验里,每次跑批自动检查特征表的max(dt)是否早于表现期起点,不满足就任务失败。
5.2 HDFS小文件多到NameNode告警
现象:跑一个简单任务,日志卡在“Listing leaf dirs”阶段很久,Yarn上申请的Container数量多到异常,NameNode内存持续上涨。原因:上游Kafka落盘每5分钟生成一个小文件,Spark写结果时默认分区数又没控制,一个分区一个小文件,三千个分区就是三千个文件。解决:ODS层按天分区,写数据前用coalesce(n)把输出文件数控制在和分区数匹配的范围;定期跑合并任务,把前一天的小文件通过INSERT OVERWRITE重新写成大文件。coalesce只减少分区不触发Shuffle,比repartition代价小得多。
5.3 Spark数据倾斜导致单Task跑半小时
现象:同一个Stage里99个任务10秒跑完,剩下1个任务跑了30分钟,整个任务卡在最后一步。原因:join或groupBy的key分布极不均衡,比如某个头部渠道的申请量占了总量的四成,这个key对应的分区负载巨大。先在Spark UI上看每个Stage的Task Duration分布,基本一眼就能定位。解决手段用加盐:
# 加盐打散热点key from pyspark.sql import functions as F df_salted = df.withColumn( "salt", F.concat(F.col("channel_id"), F.lit("_"), (F.rand() * 100).cast("int")) ) df_agg = df_salted.groupBy("salt", "user_id").agg(...)加盐后热点key先拆成100个随机后缀并行聚合,再去掉salt做最终聚合。要注意加盐会改变聚合语义,如果聚合结果需要精确按原key输出,最后一步必须再按原key聚合一次。
5.4 Executor内存与GC导致任务反复失败
现象:任务在Shuffle阶段反复报OutOfDirectMemoryError或ExecutorLostFailure,看GC日志发现Full GC频率极高。原因:executor内存设得太大,比如一台64G内存的节点只跑2个executor、每个32G,堆内GC停顿动辄几十秒,Yarn误判executor失联并杀掉。解决:单个executor内存控制在12到16g以内,宁可增加executor数量也不要堆大内存。spark.memory.fraction默认0.6,如果任务以Shuffle为主,可以调到0.7给执行内存更多空间;spark.memory.storageFraction默认0.5,缓存的数据量大时适当调低,避免RDD缓存挤占Shuffle内存。这些参数没有绝对最优,跑一轮任务看Spark UI里Execution和Storage的内存曲线,哪个接近上限就调哪个。
5.5 Hadoop与Spark版本兼容性
现象:打包好的Jar提交到集群,启动即报java.lang.NoSuchMethodError或ClassNotFoundException,看堆栈指向Hadoop的某个类。原因:本地编译用的Hadoop版本和集群不一致。比如本地用Hadoop 3.3编译,集群还是Hadoop 2.7,运行时调用新版本才有的方法直接崩。解决:pom里Hadoop client的版本号改成集群实际版本;用官方预编译版本时,认准Spark发行包名称里标注的Hadoop版本;如非必要不要自己编译Spark源码,这是无数人踩出来的血泪经验。
6. 验证模型与滚动更新:一个跑批工程师的实务技巧
6.1 按时间切片而不是随机切分
训练集和测试集的划分,信贷场景必须按时间切,不能随机切。原因是客户行为随时间漂移,随机切分会让训练集和测试集共享同一时间段的信息,离线AUC虚高。我习惯用前12个月做训练集,后3个月做验证集,最后1个月做测试集。
评估指标除了AUC,还要看KS和PR曲线。KS表示好坏客户累计分布的最大差距,风控里阈值通常取KS最大处;PR曲线关注坏客户召回率,比AUC更贴近业务损失。评估代码可以这样写:
from pyspark.ml.evaluation import BinaryClassificationEvaluator evaluator = BinaryClassificationEvaluator(labelCol="label", metricName="areaUnderROC") auc = evaluator.evaluate(predictions) print(f"验证集AUC: {auc:.4f}")6.2 阈值校准与滚动更新的个人习惯
阈值校准是我每次上线前必做的一步。如果训练时做过下采样,模型输出的概率不等于真实违约概率,需要在原始分布上重新校准:取原始样本中预测概率的分位数,找到整体坏账率对应的阈值。这个操作在风控里叫截断阈值校准,不校准就上线,模型分数分布会整体偏移。
滚动更新方面,我的习惯是每次跑批只保留最近12个月的训练样本,删除更早的数据,Spark里按日期分区过滤即可。每次调参之前先跑一版当前参数作为baseline,把AUC、KS、每个特征的系数分布存进ADS层,方便对比。模型的版本号跟着跑批日期走,评分结果表里同时记录model_version和run_date,这样出了问题能回溯到具体某一版模型。这里面的教训是:先搭好数据和评估流水的链路,再动模型,否则调参只是在自欺欺人。希望这些踩坑经验能帮你少走一段弯路。
本文还有配套的精品资源,点击获取