简介:这份资源是面向大数据与金融风控方向开发者、学生及求职者的完整项目源码,基于Hadoop与Spark构建金融信贷风险评估与管理系统,可用于课程设计、毕业项目或技术练手。压缩包共69个文件,约72KB,以36个Java与8个Scala文件承载核心业务与计算逻辑,12个XML及5个properties负责配置与依赖管理,另含SQL建表脚本、JSON数据、前端JS与说明文档,覆盖数据摄入、预处理、模型训练、风险评估到可视化各环节。项目结合Hadoop批量处理历史借贷数据与Spark实时分析申请流,涉及逻辑回归、决策树、随机森林等算法预测违约概率,并包含数据集成、安全隐私与监控运维等模块。目前已有1219人学习下载,适合希望深入理解大数据技术在金融信贷风控中落地实践的读者参考与二次开发。
1. 从一份信贷风控源码说起:Hadoop 和 Spark 到底在系统里扛了什么活
金融信贷风控这个场景,数据量不大不小,但结构特别拧巴。一边是几百 GB 的历史借贷流水、还款记录、征信快照,另一边是每秒都在刷新的申请事件、设备指纹、行为埋点。很多团队一开始用 MySQL 硬扛,跑批跑到凌晨三点,特征算不完,模型上不了线。这时候 Hadoop 和 Spark 的组合就出现了——Hadoop 负责把海量冷数据稳稳存住,Spark 负责把特征工程和评分卡计算压进可接受的时间窗口。
这份「基于 Hadoop、Spark 的大数据金融信贷风控系统源码」,本质上是一套离线为主、准实时为辅的风控数据流水线。它解决的核心问题是:把分散在业务库、日志、第三方征信里的原始数据,经过清洗、关联、聚合,变成模型能直接吃的宽表,再跑出每个借款人的风险分。适合谁看?做大数据毕业设计的学生、刚转风控数据开发的工程师、以及想从零搭一套可演示风控链路的团队。源码不是拿来直接上生产的,但它的分层思路和参数配置,能让你少走很多弯路。
2. 环境先跑通:Hadoop 伪分布式和 Spark 本地模式的落地步骤
2.1 为什么先搭伪分布式而不是直接上集群
很多人一上来就想搞三节点集群,结果卡在 SSH 免密和网络配置上,三天没跑通一个 WordCount。我的血泪经验是:先用伪分布式把 HDFS 和 YARN 的交互逻辑摸清楚,再横向扩展。伪分布式能让你看到 NameNode、DataNode、ResourceManager、NodeManager 四个进程同时起来,数据块怎么切、副本怎么放、任务怎么调度,这些在单机上都能观察。等这套跑顺了,集群搭建就是复制配置文件、改主机名的事。
Hadoop 官网的安装文档其实写得够细,但新手容易忽略的是 JDK 版本和系统用户权限。我一般用 JDK 8,建一个独立的 hadoop 用户,所有目录都放在这个用户下,避免 root 和普通用户混用导致权限报错。
2.2 从零开始安装 Hadoop 的关键命令与配置
先确认 Java 环境,然后下载 Hadoop 二进制包解压。下面这套命令是我在 Ubuntu 上反复用过的,路径按自己习惯改。
# 创建 hadoop 用户并设置密码 sudo useradd -m hadoop -s /bin/bash sudo passwd hadoop sudo adduser hadoop sudo # 切换到 hadoop 用户 su - hadoop # 下载并解压 Hadoop(版本按官网最新稳定版选,这里以 3.x 为例) wget https://downloads.apache.org/hadoop/common/hadoop-3.3.6/hadoop-3.3.6.tar.gz tar -zxvf hadoop-3.3.6.tar.gz -C /home/hadoop/ mv /home/hadoop/hadoop-3.3.6 /home/hadoop/hadoop # 配置环境变量 echo 'export HADOOP_HOME=/home/hadoop/hadoop' >> ~/.bashrc echo 'export PATH=$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin' >> ~/.bashrc source ~/.bashrc逻辑说明:这段命令做了三件事——建独立用户、解压安装包、配环境变量。参数上注意-C指定解压目录,避免散落在当前目录。环境变量里HADOOP_HOME是后续所有配置文件的基准路径,PATH加上 bin 和 sbin 才能直接敲hdfs、start-dfs.sh这些命令。
接下来改四个核心配置文件,都在$HADOOP_HOME/etc/hadoop/下。
<!-- core-site.xml --> <configuration> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:9000</value> </property> <property> <name>hadoop.tmp.dir</name> <value>/home/hadoop/hadoop/tmp</value> </property> </configuration> <!-- hdfs-site.xml --> <configuration> <property> <name>dfs.replication</name> <value>1</value> </property> <property> <name>dfs.namenode.name.dir</name> <value>/home/hadoop/hadoop/data/namenode</value> </property> <property> <name>dfs.datanode.data.dir</name> <value>/home/hadoop/hadoop/data/datanode</value> </property> </configuration>参数说明:fs.defaultFS指定 HDFS 的访问地址,伪分布式下就是 localhost 加端口。dfs.replication设成 1,因为只有一个 DataNode,设 3 会一直报副本不足。hadoop.tmp.dir和两个 data 目录必须手动创建,否则启动时直接抛异常。
YARN 的配置稍微多一点,重点是yarn-site.xml里的yarn.nodemanager.aux-services要设成mapreduce_shuffle,不然 MapReduce 任务跑不起来。内存参数根据自己机器来,8G 内存的机器给 NodeManager 分 4G 比较稳。
# 格式化 NameNode(只执行一次) hdfs namenode -format # 启动 HDFS 和 YARN start-dfs.sh start-yarn.sh # 验证进程 jpsjps应该看到 NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager 五个进程。少一个就去logs目录翻日志,九成是配置文件里路径没建或者端口被占。
2.3 Spark 本地模式跑通信贷风控的第一个特征计算
Hadoop 跑通后,Spark 的安装反而简单。下载预编译版,解压,配一个SPARK_HOME就行。本地模式不需要改集群配置,直接spark-shell或pyspark就能交互。
# 解压 Spark tar -zxvf spark-3.5.0-bin-hadoop3.tgz -C /home/hadoop/ mv /home/hadoop/spark-3.5.0-bin-hadoop3 /home/hadoop/spark # 环境变量 echo 'export SPARK_HOME=/home/hadoop/spark' >> ~/.bashrc echo 'export PATH=$PATH:$SPARK_HOME/bin' >> ~/.bashrc source ~/.bashrc # 启动 pyspark pyspark --master local[4] --executor-memory 2g--master local[4]表示用 4 个本地线程模拟并行,--executor-memory控制每个执行器的内存。本地模式下这个参数影响不大,但养成习惯,上集群时就知道怎么调。
下面这段 PySpark 代码模拟了信贷风控里最常见的「用户借款行为聚合」——从原始流水里算出每个用户的历史借款次数、平均借款金额、最大逾期天数。
from pyspark.sql import SparkSession from pyspark.sql.functions import col, count, avg, max, when # 初始化 SparkSession,指定应用名便于在 YARN 上追踪 spark = SparkSession.builder \ .appName("CreditRiskFeature") \ .config("spark.sql.shuffle.partitions", "200") \ .getOrCreate() # 读取 HDFS 上的原始借贷流水(Parquet 格式,列式存储省 IO) df = spark.read.parquet("hdfs://localhost:9000/data/loan_records/") # 按用户聚合:借款次数、平均金额、最大逾期天数 user_features = df.groupBy("user_id").agg( count("loan_id").alias("loan_count"), avg("loan_amount").alias("avg_loan_amount"), max("overdue_days").alias("max_overdue_days") ) # 衍生一个风险标记:逾期超过 90 天标记为高风险 user_features = user_features.withColumn( "high_risk_flag", when(col("max_overdue_days") > 90, 1).otherwise(0) ) # 写回 HDFS,供后续模型训练使用 user_features.write.mode("overwrite").parquet( "hdfs://localhost:9000/data/user_features/" ) spark.stop()逻辑说明:先读 Parquet 格式的原始流水,Parquet 比 CSV 省空间且带 schema,适合风控这种字段多的场景。groupBy加agg是 Spark SQL 里最核心的聚合模式,count、avg、max分别对应借款频次、金额均值、最坏逾期。withColumn加when做条件衍生,这是评分卡里常用的二值化手段。最后write.mode("overwrite")覆盖写入,避免重复跑导致数据翻倍。
参数上重点说spark.sql.shuffle.partitions,默认 200。数据量小的时候 200 个分区会导致大量小任务,调度开销比计算还大。我一般按数据量估:每 100MB 数据给 1 到 2 个分区。本地测试可以设成 8 或 16,上集群再调大。
3. 风控核心链路:从原始流水到风险评分的 Spark 实现
3.1 数据分层建模:ODS、DWD、DWS 在信贷场景怎么切
风控系统的数据不能一锅炖。我习惯按三层走:ODS 层放原始数据,不做任何清洗,保留所有字段和脏数据;DWD 层做清洗和标准化,比如把不同来源的日期格式统一、把金额单位统一成元、把用户 ID 做脱敏映射;DWS 层做聚合宽表,就是模型直接用的特征表。
为什么这么切?因为风控模型迭代快,今天用 A 特征,明天可能换 B 特征。如果直接从原始数据算,每次都要重跑全量清洗,耗时且容易出错。分层之后,DWD 层是稳定的,DWS 层可以按需重建。源码里通常会有ods_loan、dwd_loan_clean、dws_user_feature这样的表名,看到这种命名就知道作者是按数仓思路做的。
在 Spark 里实现分层,可以用不同的数据库或不同的 HDFS 目录来隔离。我一般用 Hive 管理元数据,Spark 通过enableHiveSupport()直接读写 Hive 表。
# 启用 Hive 支持,方便用 SQL 管理分层表 spark = SparkSession.builder \ .appName("CreditRiskLayer") \ .enableHiveSupport() \ .getOrCreate() # ODS 层:原始数据直接映射,不做转换 spark.sql(""" CREATE TABLE IF NOT EXISTS ods_loan ( loan_id STRING, user_id STRING, loan_amount DOUBLE, loan_date STRING, overdue_days INT, source STRING ) STORED AS PARQUET """) # DWD 层:清洗日期格式、过滤金额为负的异常记录 spark.sql(""" CREATE TABLE IF NOT EXISTS dwd_loan_clean AS SELECT loan_id, user_id, loan_amount, to_date(loan_date, 'yyyy-MM-dd') AS loan_date, overdue_days, source FROM ods_loan WHERE loan_amount > 0 """) # DWS 层:按用户聚合特征 spark.sql(""" CREATE TABLE IF NOT EXISTS dws_user_feature AS SELECT user_id, COUNT(loan_id) AS loan_count, AVG(loan_amount) AS avg_loan_amount, MAX(overdue_days) AS max_overdue_days, SUM(CASE WHEN overdue_days > 90 THEN 1 ELSE 0 END) AS high_risk_count FROM dwd_loan_clean GROUP BY user_id """)逻辑说明:ODS 建表时字段类型要跟原始数据对齐,loan_date先存 STRING,到 DWD 再转 DATE。DWD 的WHERE loan_amount > 0过滤掉冲正或测试数据,这是风控里常见的脏数据。DWS 用CASE WHEN做条件计数,比先过滤再 join 效率高。
参数上注意STORED AS PARQUET,风控宽表字段多,Parquet 的列裁剪能大幅减少 IO。如果数据量到 TB 级,可以再加分区,比如按loan_date的月份分区。
3.2 用 Spark SQL 算逾期率和 WOE 分箱
风控模型里离不开逾期率和 WOE(Weight of Evidence)。逾期率是基础指标,WOE 是评分卡特征筛选的核心。Spark SQL 算这两个指标比写 MapReduce 快得多,而且代码量少。
先算逾期率。假设 DWD 层有一张dwd_loan_clean表,包含loan_date、overdue_days、loan_amount。
-- 按月统计逾期率:逾期天数大于 0 算逾期 SELECT DATE_FORMAT(loan_date, 'yyyy-MM') AS loan_month, COUNT(loan_id) AS total_loans, SUM(CASE WHEN overdue_days > 0 THEN 1 ELSE 0 END) AS overdue_loans, ROUND(SUM(CASE WHEN overdue_days > 0 THEN 1 ELSE 0 END) / COUNT(loan_id), 4) AS overdue_rate FROM dwd_loan_clean GROUP BY DATE_FORMAT(loan_date, 'yyyy-MM') ORDER BY loan_month;逻辑说明:DATE_FORMAT把日期截到月,CASE WHEN做条件计数,最后除出比率。ROUND保留四位小数,方便看趋势。这个查询在 Spark SQL 里会走一次 shuffle,数据量大时可以把结果写回 Hive 表,避免重复计算。
WOE 的计算稍微复杂,需要先分箱。分箱就是把连续变量切成若干区间,比如借款金额分成 0-1万、1-5万、5-10万、10万以上。Spark SQL 里可以用CASE WHEN做等频或等距分箱。
-- 对借款金额做等距分箱,并计算每个箱的 WOE WITH binned AS ( SELECT loan_id, CASE WHEN loan_amount < 10000 THEN '0-1w' WHEN loan_amount < 50000 THEN '1-5w' WHEN loan_amount < 100000 THEN '5-10w' ELSE '10w+' END AS amount_bin, CASE WHEN overdue_days > 0 THEN 1 ELSE 0 END AS is_overdue FROM dwd_loan_clean ), bin_stats AS ( SELECT amount_bin, COUNT(*) AS total, SUM(is_overdue) AS bad, COUNT(*) - SUM(is_overdue) AS good FROM binned GROUP BY amount_bin ), total_stats AS ( SELECT SUM(bad) AS total_bad, SUM(good) AS total_good FROM bin_stats ) SELECT b.amount_bin, b.total, b.bad, b.good, ROUND(LN((b.bad / t.total_bad) / (b.good / t.total_good)), 4) AS woe FROM bin_stats b CROSS JOIN total_stats t ORDER BY b.amount_bin;逻辑说明:CTE 先分箱,再按箱统计好坏样本数,最后用LN算 WOE。WOE 的公式是ln(坏样本占比 / 好样本占比),值越大说明这个箱的风险越高。CROSS JOIN把总好坏样本数拼到每一行,避免写子查询。
参数上注意分箱边界。等距分箱简单但可能把风险差异大的区间合并,等频分箱更均匀但边界不好解释。我一般先用等距跑一版看 WOE 单调性,不单调再调边界。WOE 的绝对值大于 0.5 就说明这个箱的区分度不错。
3.3 把特征表喂给逻辑回归:Spark MLlib 训练与评估
特征算完,下一步是训练模型。风控里逻辑回归是主力,因为可解释性强,监管能看懂。Spark MLlib 的LogisticRegression支持弹性网络正则,适合处理特征共线性。
from pyspark.ml.feature import VectorAssembler, StandardScaler from pyspark.ml.classification import LogisticRegression from pyspark.ml.evaluation import BinaryClassificationEvaluator # 读取 DWS 特征表 feature_df = spark.sql("SELECT * FROM dws_user_feature") # 组装特征向量,排除 user_id 和标签列 assembler = VectorAssembler( inputCols=["loan_count", "avg_loan_amount", "max_overdue_days", "high_risk_count"], outputCol="raw_features" ) assembled = assembler.transform(feature_df) # 标准化:逻辑回归对量纲敏感,金额和次数差几个数量级 scaler = StandardScaler( inputCol="raw_features", outputCol="features", withMean=True, withStd=True ) scaled = scaler.fit(assembled).transform(assembled) # 划分训练集和测试集 train, test = scaled.randomSplit([0.8, 0.2], seed=42) # 训练逻辑回归,设置正则化参数 lr = LogisticRegression( featuresCol="features", labelCol="high_risk_flag", maxIter=100, regParam=0.01, elasticNetParam=0.5 ) model = lr.fit(train) # 预测并评估 AUC predictions = model.transform(test) evaluator = BinaryClassificationEvaluator( labelCol="high_risk_flag", rawPredictionCol="rawPrediction", metricName="areaUnderROC" ) auc = evaluator.evaluate(predictions) print(f"AUC: {auc:.4f}")逻辑说明:VectorAssembler把多个特征列拼成一个向量,这是 MLlib 的输入格式要求。StandardScaler做 Z-score 标准化,withMean和withStd都设 True。randomSplit按 8:2 分,seed固定保证可复现。LogisticRegression的regParam控制正则强度,elasticNetParam设 0.5 表示 L1 和 L2 各占一半。BinaryClassificationEvaluator算 AUC,风控里 AUC 到 0.75 以上算可用。
参数上重点说regParam。太大模型欠拟合,太小过拟合。我一般从 0.01 开始,看训练集和测试集的 AUC 差距,差距大就调大正则。maxIter默认 100,一般够用,不收敛再加到 200。
4. 避坑与排查:Hadoop 和 Spark 在风控场景里最容易翻车的地方
4.1 坑一:HDFS 小文件太多导致 NameNode 内存爆
现象:跑完特征计算后,HDFS 上生成了几万个几十 KB 的小文件,NameNode 内存持续上涨,最后 Web UI 打不开。
原因:Spark 的write默认按分区数输出文件,如果分区设成 200,每个分区写一个文件,数据量小的时候就是 200 个小文件。风控系统每天跑批,日积月累小文件数量爆炸。
解决:写入前用coalesce或repartition合并分区。coalesce不触发 shuffle,适合减少分区;repartition触发 shuffle,适合增加分区。我一般在写入前加一句df.coalesce(10).write...,把文件数控制在 10 个以内。另外可以开 Hive 的合并小文件参数,但治标不治本,源头控制更靠谱。
4.2 坑二:Spark 内存溢出,Executor 频繁 GC
现象:任务跑到某个 stage 卡住,日志里全是GC overhead limit exceeded,或者直接OutOfMemoryError。
原因:风控宽表字段多,groupBy和join时数据倾斜,某个 key 的数据量特别大,一个 Executor 扛不住。另外spark.executor.memory设得太小,或者spark.sql.shuffle.partitions太大导致每个分区数据量不均。
解决:先看 Spark UI 的 stage 详情,找数据量最大的 task。如果是倾斜,把大 key 加随机前缀打散,算完再去前缀。如果是内存不够,调大spark.executor.memory和spark.executor.memoryOverhead。我一般把memoryOverhead设成 executor 内存的 10% 到 15%,给 JVM 堆外留空间。另外spark.sql.adaptive.enabled设成 true,让 Spark 自动处理倾斜。
4.3 坑三:Hadoop 和 Spark 版本不匹配导致 NoSuchMethodError
现象:Spark 任务提交到 YARN 后报NoSuchMethodError或ClassNotFoundException,指向 Hadoop 的某个类。
原因:Spark 预编译版绑定了特定 Hadoop 版本,比如spark-3.5.0-bin-hadoop3对应 Hadoop 3.x。如果集群装的是 Hadoop 2.x,类路径里的 jar 包版本冲突。
解决:下载 Spark 时看清后缀,bin-hadoop3对应 Hadoop 3,bin-hadoop2.7对应 Hadoop 2.7。如果已经装错,可以重新下载匹配版本,或者手动替换SPARK_HOME/jars下的 Hadoop 相关 jar。我一般直接重下,省得排查依赖。
4.4 坑四:YARN 队列资源不足,任务一直 pending
现象:Spark 任务提交后一直处于ACCEPTED状态,不进入RUNNING。
原因:YARN 的队列资源被占满,或者spark.executor.instances设得太大,超过了队列容量。风控跑批经常和别的任务抢资源。
解决:先看 YARN 的 ResourceManager Web UI,确认队列剩余资源。调小spark.executor.instances或spark.executor.cores,让任务能挤进去。如果队列是 Capacity Scheduler,可以临时调大队列容量,但生产环境要走审批。我一般把跑批任务设成低优先级队列,避免和实时任务抢。
4.5 坑五:日期格式不统一导致特征计算全错
现象:逾期率算出来是 0 或者异常高,检查发现loan_date有的存2024-01-01,有的存2024/01/01,还有的存时间戳。
原因:不同业务系统写入 ODS 层时没有统一格式,Spark 的to_date遇到不匹配的格式返回 null,后续聚合全乱。
解决:在 DWD 层做强制转换,用CASE WHEN判断格式再转。或者用to_date的多种格式参数,但 Spark 的to_date只支持一种格式。我一般写一个 UDF 做多格式解析,或者直接在 ODS 接入时用 Flume 或 Kafka Connect 做预处理。风控数据里日期是核心字段,格式不统一是最大的玄学问题。
5. 进阶技巧:用 Spark 内存监测和分区调优把跑批时间压下来
跑通链路只是第一步,真正上生产要关注跑批时间。我拿一个实际案例说:某信贷风控系统,原始流水 500GB,特征表 200 个字段,最初跑批要 4 小时,调优后压到 50 分钟。核心做了三件事。
第一,开启动态资源分配。Spark 默认会一次性申请所有 Executor,跑批任务前期数据量小,后期数据量大,资源浪费严重。开启动态分配后,Spark 根据 stage 的数据量自动增减 Executor。
spark-submit \ --master yarn \ --deploy-mode cluster \ --conf spark.dynamicAllocation.enabled=true \ --conf spark.dynamicAllocation.minExecutors=5 \ --conf spark.dynamicAllocation.maxExecutors=50 \ --conf spark.dynamicAllocation.initialExecutors=10 \ --conf spark.dynamicAllocation.executorIdleTimeout=60s \ --conf spark.shuffle.service.enabled=true \ --conf spark.sql.adaptive.enabled=true \ --conf spark.sql.adaptive.coalescePartitions.enabled=true \ --conf spark.sql.adaptive.skewJoin.enabled=true \ --executor-memory 8g \ --executor-cores 4 \ --num-executors 10 \ credit_risk_feature.py参数说明:minExecutors和maxExecutors控制弹性范围,executorIdleTimeout设 60 秒,空闲一分钟就释放。spark.shuffle.service.enabled必须开,不然动态分配释放 Executor 时 shuffle 数据会丢。spark.sql.adaptive.enabled是 Spark 3.x 的自适应查询执行,能自动合并小分区、处理倾斜 join。coalescePartitions和skewJoin是它的子功能,建议都开。
第二,用 Spark UI 的内存监测定位 GC 瓶颈。Spark UI 的 Executors 页面能看到每个 Executor 的 GC 时间占比。如果某个 Executor 的 GC 时间超过 10%,说明内存不够或者对象太多。我一般先看 Storage 页面,确认有没有缓存了不该缓存的数据。风控特征表通常不需要cache(),因为只用一次。如果确实要缓存,用MEMORY_AND_DISK级别,别用MEMORY_ONLY。
第三,分区调优。spark.sql.shuffle.partitions默认 200,对于 500GB 数据,每个分区 2.5GB,太大。我按每个分区 128MB 到 256MB 来估,500GB 大概需要 2000 到 4000 个分区。但分区太多调度开销大,所以配合自适应执行,让 Spark 自动合并。
# 在代码里动态设置 shuffle 分区数 spark.conf.set("spark.sql.shuffle.partitions", "2000") # 对大表 join 前先做分桶,避免 shuffle spark.sql(""" CREATE TABLE dwd_loan_bucketed USING PARQUET CLUSTERED BY (user_id) INTO 256 BUCKETS AS SELECT * FROM dwd_loan_clean """)分桶是比分区更细的粒度,CLUSTERED BY (user_id) INTO 256 BUCKETS把相同 user_id 的数据写到同一个桶文件。后续 join 时如果两边都按 user_id 分桶,Spark 可以直接做 bucket join,省掉 shuffle。风控里用户 ID 是核心关联键,分桶收益很大。
最后说一个验证方法:跑批前后对比spark.sql.adaptive开和关的耗时。我一般跑两次,一次关自适应,一次开,看 stage 数量和 shuffle 数据量。如果开自适应后 stage 数减少 30% 以上,说明调优有效。另外看 YARN 的 ResourceManager 日志,确认没有容器被 kill。
这些调优手段不是一次全上,先开自适应执行,再调分区数,最后上分桶。每改一个参数跑一次,记录耗时,找到瓶颈再动下一个。我踩过的最大坑是一次性改了五个参数,结果跑批时间反而变长,排查了两天才发现是executor-cores设成 8 导致 HDFS 并发读写太高,磁盘 IO 打满。后来改成 4 就正常了。调优是个细活,别贪多。
希望帮到你。
本文还有配套的精品资源,点击获取