简介这份资源是面向大数据与金融科技方向开发者、学生及求职者的完整项目源码聚焦于利用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、node3node1 兼做 NameNode 和 ResourceManager。配置文件关键参数建议值说明core-site.xmlfs.defaultFShdfs://node1:9000NameNode 地址core-site.xmlhadoop.tmp.dir/data/hadoop/tmp别用默认 /tmp重启丢数据hdfs-site.xmldfs.replication3三节点就设 3hdfs-site.xmldfs.namenode.name.dir/data/hadoop/nn元数据目录hdfs-site.xmldfs.datanode.data.dir/data/hadoop/dn数据块目录yarn-site.xmlyarn.nodemanager.resource.memory-mb8192按物理内存的 70% 给yarn-site.xmlyarn.scheduler.maximum-allocation-mb4096单个容器上限mapred-site.xmlmapreduce.framework.nameyarn走 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 -formatNameNode 的 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条件永远不成立。解决读进来先castF.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里建否则连接数爆炸。推之前把特征序列化成 JSONkey 用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 的版本兼容性也要盯紧源码里如果写死了某个版本换版本前先看官方兼容矩阵。希望帮到你。本文还有配套的精品资源点击获取