简介:这份资源是面向大数据与金融科技方向开发者、学生及求职者的完整项目源码,聚焦于利用Hadoop与Spark构建金融信贷风险评估与管理系统,帮助读者理解大数据技术在信贷风控场景中的落地方式。压缩包共69个文件,约72KB,以36个Java文件与8个Scala文件构成核心业务与计算逻辑,辅以12个XML配置、5个properties参数文件及SQL、JSON、JS等资源,覆盖数据摄入、预处理、模型训练、风险评估与可视化等环节。项目结合Hadoop的HDFS与MapReduce处理海量历史借贷数据,并借助Spark内存计算实现实时风险评分,涉及逻辑回归、决策树、随机森林等机器学习算法,同时体现批处理与流处理结合、多源数据集成、安全隐私保护及可扩展性等设计思路。目前已有1219人学习下载,适合希望深入实践大数据金融风控的开发者参考与二次开发。
1. 从一份信贷风控源码说起:Hadoop 和 Spark 到底在系统里扛了什么活
金融信贷风控这个场景,数据量不大不小,但结构特别拧巴。一边是借款人基本信息、合同、还款计划这类规整的业务表,另一边是设备指纹、操作日志、第三方多头借贷查询记录这类半结构化甚至非结构化的数据。单机 MySQL 跑到几十万笔借据就开始喘,更别说做变量衍生和模型回溯。所以当有人把 Hadoop 和 Spark 塞进一个信贷风控系统里,我第一反应不是"炫技",而是这个组合确实对得上需求:Hadoop 负责把海量原始数据低成本地存下来,Spark 负责把特征工程和批量评分跑快。
这份源码标题里"基于 Hadoop、Spark 的大数据金融信贷风险控系统",拆开看就是三层:底层是 HDFS 存原始借据、还款流水、征信报文;中间是 Spark 做特征加工和风险指标计算;上层是风控规则引擎和评分卡输出。适合谁看?做大数据毕业设计的同学能拿到一套完整链路,做信贷系统开发的工程师能看清离线特征怎么落到线上决策,做数据平台的同学能对照自己的集群规划。接下来我不讲空概念,直接按"集群怎么搭、数据怎么进、特征怎么算、坑在哪"这条线走一遍。
2. 集群底座:Hadoop 伪分布式到三节点,先把存储和调度立住
2.1 为什么风控系统先要 HDFS 而不是直接上对象存储
信贷风控的数据有个特点:写一次、读很多次,而且读的时候往往是全量扫描做回溯。比如你要验证一个新规则在过去 12 个月的表现,就得把历史借据全捞出来重算。这种访问模式下,HDFS 的机架感知副本机制比对象存储的按请求计费更可控,尤其是自建集群时,三副本带来的容错是实打实的。
另一个原因是生态。Spark 读 HDFS 是原生接口,不用额外适配层。源码里如果用了 Hive 做元数据管理,那 HDFS 更是绕不开的底座。常见做法是 NameNode 单独一台,DataNode 和 NodeManager 混布,小集群三台起步。
2.2 三节点 Hadoop 集群的最小配置清单
先给一份我常用的配置表,按这个改完基本能跑起来。主机名假设是 node1、node2、node3,node1 兼做 NameNode 和 ResourceManager。
| 配置文件 | 关键参数 | 建议值 | 说明 |
|---|---|---|---|
| core-site.xml | fs.defaultFS | hdfs://node1:9000 | NameNode 地址 |
| core-site.xml | hadoop.tmp.dir | /data/hadoop/tmp | 别用默认 /tmp,重启丢数据 |
| hdfs-site.xml | dfs.replication | 3 | 三节点就设 3 |
| hdfs-site.xml | dfs.namenode.name.dir | /data/hadoop/nn | 元数据目录 |
| hdfs-site.xml | dfs.datanode.data.dir | /data/hadoop/dn | 数据块目录 |
| yarn-site.xml | yarn.nodemanager.resource.memory-mb | 8192 | 按物理内存的 70% 给 |
| yarn-site.xml | yarn.scheduler.maximum-allocation-mb | 4096 | 单个容器上限 |
| mapred-site.xml | mapreduce.framework.name | yarn | 走 YARN 调度 |
配置改完,格式化和启动的命令如下:
# 在 node1 上格式化 NameNode,只做一次 hdfs namenode -format # 启动 HDFS start-dfs.sh # 启动 YARN start-yarn.sh # 验证进程,node1 应该有 NameNode、ResourceManager、DataNode、NodeManager jps # 建风控系统的数据目录 hdfs dfs -mkdir -p /risk/raw/loan hdfs dfs -mkdir -p /risk/raw/repay hdfs dfs -mkdir -p /risk/warehouse逻辑说明:hdfs namenode -format会清空元数据目录,重复执行会导致集群 ID 不一致,DataNode 起不来,这是新手最常见的翻车点。jps是排查进程是否齐全的第一手段,少一个进程就去对应日志目录翻。目录规划上,raw 放原始数据,warehouse 放 Hive 表,后面 Spark 写特征也往 warehouse 下挂。
参数说明:dfs.replication设 3 是因为三节点刚好每个节点一份,容错和空间平衡。yarn.nodemanager.resource.memory-mb别设成物理内存全量,操作系统和 DataNode 自己还要吃内存,留 30% 是血泪经验。
2.3 把原始借据数据灌进 HDFS 的两种方式
源码里通常会给一份 CSV 或 SQL 导出文件。小数据量直接 put,大数据量走 Sqoop 或 DataX。我一般先用 put 验证链路:
# 本地文件上传到 HDFS hdfs dfs -put ./loan_2023.csv /risk/raw/loan/ # 查看文件是否完整 hdfs dfs -ls /risk/raw/loan/ hdfs dfs -cat /risk/raw/loan/loan_2023.csv | head -5如果数据在 MySQL 里,用 Sqoop 抽:
sqoop import \ --connect jdbc:mysql://dbserver:3306/credit \ --username risk_reader \ --password-file /risk/.mysql.pwd \ --table loan_contract \ --target-dir /risk/raw/loan \ --fields-terminated-by '\001' \ --num-mappers 4逻辑说明:--fields-terminated-by '\001'用不可见字符做分隔,避免业务字段里出现逗号导致列错位,这是金融数据里特别容易踩的坑,因为备注字段什么都可能写。--num-mappers 4控制并行度,别设太大,MySQL 那边连接数扛不住。
3. Spark 特征工程:把借据流水算成风控变量
3.1 风控变量为什么必须在 Spark 里做而不是 SQL
信贷风控的核心变量,比如"近 3 个月申请次数""当前逾期金额占比""最大连续逾期天数",都是窗口函数加聚合。MySQL 8 也能写窗口函数,但数据量上到千万级借据、上亿条还款流水时,单机 SQL 就跑不动了。Spark 的优势在于把窗口计算分布式化,而且能直接读 HDFS 上的原始文件,省掉导入导出。
另一个原因是特征回溯。模型上线后要监控变量稳定性,得按历史时点重算变量。Spark 的 DataFrame API 写这种时点回溯逻辑比 SQL 清晰,尤其是配合Window.partitionBy().orderBy()的时候。
3.2 用 PySpark 算三个典型风控变量
下面这段代码算三个变量:申请次数、逾期金额占比、最大连续逾期天数。假设原始数据已经以 Parquet 格式存在 HDFS 上。
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window spark = SparkSession.builder \ .appName("risk_feature") \ .config("spark.sql.shuffle.partitions", "200") \ .enableHiveSupport() \ .getOrCreate() # 读借据表 loan = spark.read.parquet("hdfs://node1:9000/risk/raw/loan") # 读还款流水 repay = spark.read.parquet("hdfs://node1:9000/risk/raw/repay") # 变量1:每个客户近3个月申请次数 apply_cnt = loan.filter( F.col("apply_time") >= F.date_sub(F.current_date(), 90) ).groupBy("cust_id").agg( F.count("loan_id").alias("apply_cnt_3m") ) # 变量2:当前逾期金额占比 overdue_ratio = repay.groupBy("cust_id").agg( F.sum(F.when(F.col("overdue_days") > 0, F.col("due_amount")).otherwise(0)).alias("overdue_amt"), F.sum("due_amount").alias("total_amt") ).withColumn( "overdue_ratio", F.col("overdue_amt") / F.col("total_amt") ).select("cust_id", "overdue_ratio") # 变量3:最大连续逾期天数,用窗口函数标记连续段 w = Window.partitionBy("cust_id").orderBy("repay_date") repay_flag = repay.withColumn( "is_overdue", F.when(F.col("overdue_days") > 0, 1).otherwise(0) ).withColumn( "grp", F.sum(F.when(F.col("is_overdue") == 0, 1).otherwise(0)).over(w) ) max_cont = repay_flag.filter(F.col("is_overdue") == 1).groupBy("cust_id", "grp").agg( F.count("repay_date").alias("cont_days") ).groupBy("cust_id").agg( F.max("cont_days").alias("max_cont_overdue_days") ) # 合并变量写入 Hive feature = apply_cnt.join(overdue_ratio, "cust_id", "left") \ .join(max_cont, "cust_id", "left") \ .fillna(0) feature.write.mode("overwrite").saveAsTable("risk.feature_cust")逻辑说明:变量 3 的连续段标记是经典套路,用累计非逾期次数做分组键,同一组内的逾期记录就是连续的。fillna(0)处理没有逾期记录的客户,避免后续模型训练出空值。spark.sql.shuffle.partitions设 200 是经验值,太小会导致单分区数据倾斜,太大产生大量小文件。
参数说明:enableHiveSupport()让 Spark 能直接读写 Hive 表,前提是 hive-site.xml 已经放到 Spark 的 conf 目录。mode("overwrite")每次全量覆盖,生产上更稳的做法是按日期分区增量写。
3.3 特征写入 Hive 后的校验动作
写完不能直接信,得校验。我一般跑三个检查:
-- 检查记录数是否和客户数对得上 SELECT COUNT(*) FROM risk.feature_cust; -- 检查关键变量是否有异常值 SELECT MAX(overdue_ratio), MIN(overdue_ratio), AVG(apply_cnt_3m) FROM risk.feature_cust; -- 检查空值比例 SELECT SUM(CASE WHEN max_cont_overdue_days IS NULL THEN 1 ELSE 0 END) / COUNT(*) FROM risk.feature_cust;逻辑说明:overdue_ratio理论上应该在 0 到 1 之间,如果出现大于 1,说明还款流水的 due_amount 有重复累加,得回去查数据源。空值比例超过 5% 就要警惕,可能是 join 的时候客户 ID 对不上。
4. 避坑与排查:集群和 Spark 作业最容易翻车的地方
4.1 DataNode 起不来,日志报 clusterID 不一致
现象:start-dfs.sh后 jps 看不到 DataNode,日志里写Incompatible clusterIDs。
原因:重复执行了hdfs namenode -format,NameNode 的 clusterID 变了,DataNode 还记着旧的。
解决:要么把 DataNode 的 data 目录清空重新格式化,要么把 NameNode 的 clusterID 手动改成和 DataNode 一致。生产上格式化只做一次,做完立刻备份元数据目录。
4.2 Spark 作业卡在最后一个 stage 不动
现象:Web UI 上看到 199 个 task 完成了,剩 1 个跑了几十分钟。
原因:数据倾斜。某个 cust_id 的流水特别多,全分到一个分区。
解决:先看spark.sql.shuffle.partitions是不是太小,调大试试。如果是热点 key,用加盐的方式打散,比如给 cust_id 拼一个随机后缀,聚合两次。源码里如果没处理倾斜,大数据量下必翻车。
4.3 读 Hive 表报 ClassNotFoundException
现象:Spark 代码里enableHiveSupport()之后读表报找不到 Hive 的类。
原因:Spark 的 conf 目录下没有 hive-site.xml,或者 Hive 的 jar 包没进 classpath。
解决:把 Hive 的 hive-site.xml 软链到$SPARK_HOME/conf/,确保spark.sql.catalogImplementation是 hive。用spark-submit时加--jars把 Hive 的依赖带上。
4.4 特征变量算出来全是 0
现象:overdue_ratio和max_cont_overdue_days全是 0。
原因:还款流水里的overdue_days字段类型是字符串,和数字比较时隐式转换失败,when条件永远不成立。
解决:读进来先cast,F.col("overdue_days").cast("int")。金融数据从 CSV 或 MySQL 抽过来,字段类型经常是 string,这是高频坑。
4.5 YARN 容器被 kill,报超出内存
现象:Spark 作业跑一半 executor 全没了,YARN 日志写Container killed by YARN for exceeding memory limits。
原因:spark.executor.memory设得比yarn.scheduler.maximum-allocation-mb还大,或者 executor 的堆外内存没算进去。
解决:executor 内存加上spark.executor.memoryOverhead要小于容器上限。一般 executor 内存设容器上限的 80%,留 20% 给 overhead。
5. 从离线特征到风控决策:评分卡接入和增量调优
5.1 把 Spark 特征表接到评分卡引擎
特征算完只是半成品,得让风控规则用上。常见做法是 Spark 把特征表写到 HDFS 或 HBase,评分卡引擎定时拉取。如果源码里带了规则引擎模块,一般是读 Hive 表然后跑 Drools 或自研的规则解析器。我一般会加一层缓存,把客户维度的特征推到 Redis,线上决策时直接查,避免每次请求都扫 Hive。
# 特征推 Redis 的简化逻辑 feature = spark.table("risk.feature_cust") feature.foreachPartition(lambda rows: push_to_redis(rows))逻辑说明:foreachPartition每个分区建一次 Redis 连接,别在foreach里建,否则连接数爆炸。推之前把特征序列化成 JSON,key 用risk:feature:{cust_id}。
5.2 变量监控:怎么知道特征算得对不对
上线后每周跑一次变量监控,看三个指标:缺失率、PSI、命中率。PSI 超过 0.1 说明变量分布漂移,可能是数据源变了或者计算逻辑有 bug。我习惯把监控结果写回 Hive 表,用 Spark 定时任务跑,出问题发告警。
| 监控指标 | 计算方式 | 阈值 | 处理动作 |
|---|---|---|---|
| 缺失率 | 空值数 / 总数 | > 5% | 查数据源和 join 逻辑 |
| PSI | 当期分布 vs 基期分布 | > 0.1 | 排查变量逻辑变更 |
| 命中率 | 非零值数 / 总数 | 波动 > 20% | 查上游数据量 |
5.3 增量调优:从全量重算到按日分区
全量重算特征表在数据量上来后越来越慢。我一般改成按日分区,每天只算增量,历史分区不动。Spark 写的时候用partitionBy("dt"),读的时候用where dt >=过滤。这样回溯也方便,指定日期范围就行。
feature.write.mode("overwrite") \ .partitionBy("dt") \ .saveAsTable("risk.feature_cust_daily")逻辑说明:分区字段选日期,别选客户 ID,否则小文件多到 NameNode 扛不住。每天一个分区,一年 365 个目录,可控。
这套东西我从头搭过一遍,最大的教训是别一上来就追求集群规模,三台机器先把链路跑通,特征算对了再扩。Hadoop 和 Spark 的版本兼容性也要盯紧,源码里如果写死了某个版本,换版本前先看官方兼容矩阵。希望帮到你。
本文还有配套的精品资源,点击获取