简介:面向大数据与金融风控开发者的可运行源码项目,基于Hadoop与Spark构建信贷风险控制系统,结合HDFS分布式存储、MapReduce批处理、Spark内存计算与流处理技术,覆盖数据摄入、预处理、模型训练、风险评估到可视化展示的全链路流程。资源共69个文件,以36个Java、8个Scala源码为核心,承担数据清洗、风险评分、机器学习算法等主逻辑;12个XML配置文件用于Spring/MyBatis等组件装配,5个Properties与SQL脚本完成环境配置、数据库初始化。压缩包仅72KB,目录结构清晰,拆分为data-source-spark-streaming与credit-risk-control两个Maven模块,并带pom.xml与IDEA运行配置,便于直接导入二次开发。前者实现多源数据实时接入与特征加工,后者包含信用评估、欺诈检测、还款能力分析及可视化接口,完整展现Hadoop+Spark整合思路。已有1219人学习下载,适合希望快速搭建可扩展大数据风控原型、深入实践源码级方案的开发者。
1. 这是一套什么源码:Hadoop+Spark 的信贷风控系统到底长什么样
如果你是做大数据开发或者想转风控方向的人,八成会在 GitHub 或各种付费源码站刷到这类标题的压缩包:基于 Hadoop、Spark 的金融信贷风险控系统源码.zip。我头一回拿到这种包的时候,心里想的是「白嫖一套完整项目」,解压完才发现它并不像普通 Web 项目那样开箱即用。它本质是一套面向小额信贷、消费金融场景的离线风控平台:用 Hadoop HDFS 存申请数据、征信报文、行为日志,用 Spark 做批量特征计算和跑分,最后把评分结果回写业务库,给审批系统一个放款或拒贷建议。这个方向适合正在积累大数据项目经验的工程师,也适合小额贷公司内部数据团队做自研风控的起步参考。真正难住大部分人的不是算法,而是把一套 zip 里的 Hadoop、Spark、Hive、MySQL 和策略引擎串起来跑通这件事本身。
2. 拆开 ZIP 先看架构:存储、计算、规则与接口怎么分工
2.1 Hadoop 负责什么:HDFS 存储与 Hive 数仓分层
金融信贷风控系统的数据底座几乎都长一个样。业务库(MySQL / Oracle)里放着用户申请主表、借款流水、还款记录;第三方征信报告往往是文件形式,一天一推,可能是 JSON 也可能是定长文本;APP 端还有埋点行为日志。这么多来源、格式不齐的数据,不可能直接丢给规则引擎跑分,所以第一件事就是统一进 HDFS。
HDFS 在这里扮演「原始数据湖」的角色。常见做法是按日期分区落地:/data/ods/apply/20240920/放当天申请记录,/data/ods/credit_report/20240920/放征信报文。然后用 Hive 建外部表映射这些目录,做 ODS(原始数据层)到 DWD(明细层)的清洗。清洗逻辑简单说就是去重、格式统一、把嵌套 JSON 展开成宽表。
CREATE EXTERNAL TABLE dwd_apply_info( apply_id STRING, user_id STRING, apply_amount DECIMAL(12,2), apply_time STRING, channel STRING, id_card_hash STRING ) PARTITIONED BY (dt STRING) STORED AS PARQUET LOCATION '/data/dwd/apply_info'; ALTER TABLE dwd_apply_info ADD PARTITION (dt='2024-09-20');这套分层结构的好处是,规则引擎和 Spark 特征计算只需要读 DWD 层,不用天天面对源数据的脏格式。对信贷场景来说,历史数据要保留完整,因为后面做违约率回溯分析时,需要把某个月的申请人群和之后 3 个月、6 个月的逾期表现关联起来,这个关联在 ODS 原始层做最保险。
2.2 Spark 负责什么:批处理跑分、规则引擎与结果落库
Hadoop 把数据管起来了,但真正让系统「有智商」的是 Spark 这一层。我见过的信贷风控源码里,Spark 的任务基本分三类:第一类是跑特征,把用户过去 30 天申请次数、过去 90 天贷款总额、征信查询次数等衍生字段算出来;第二类是跑规则引擎,把特征值喂给规则集,输出命中结果;第三类是模型打分,加载训练好的逻辑回归或 XGBoost 模型,对全量申请批量预测违约概率。
这三类任务在代码里通常都是一个个 Spark 作业,打包成 jar 后用spark-submit提交。调度方式比较常见的做法是用 Crontab 或者 Azkaban 每天凌晨跑 T+1 批处理;如果时效要求高一些,就上 Spark Structured Streaming 做小时级甚至分钟级处理。
spark-submit \ --class com.credit.risk.engine.RuleEngineJob \ --master yarn \ --deploy-mode cluster \ --driver-memory 4g \ --executor-memory 8g \ --executor-cores 4 \ --num-executors 20 \ ./credit-risk-engine.jar \ --biz_date 2024-09-20这段命令里值得留意的是--deploy-mode cluster。很多新手在伪分布式环境跑通后用同一个参数上集群,结果日志里一直报「找不到主类」,其实是把 driver 跑在了集群节点上,日志看不到也查不到,血泪经验。后面第 5 章我会展开说这个坑。另外一个重点是 executor 数量和内存的配比,它直接决定风控批处理能不能在凌晨两小时的窗口内跑完。规则引擎这类任务以内存 shuffle 为主,executor 内存给 4g 通常不够,8g 起步在信贷数据量下是正常操作。
2.3 系统里还有哪些躲不开的模块:MySQL、Redis 与接口层
别以为用上了 Hadoop 和 Spark,MySQL 就没用了。实际源码里 MySQL 至少干三件事:存风控规则配置、存跑分结果、存审批回调状态。规则配置放在 MySQL 里而不是写死在代码里,理由是风控策略调整频率非常高,信贷产品上线后几乎每周都要调规则阈值,把这个做成配置项才能让非技术人员也能改。Redis 在这个系统里做的是批量跑分完成后的结果缓存,审批系统实时查询时先从 Redis 拿,拿不到再查 MySQL,避免每次审批都去打宽表。
接口层一般是一个 Spring Boot 服务,它做的事情是把规则引擎产出的分数和命中规则翻译成人话:拒绝、人工审核、自动通过,还要附上拒绝原因码。所以你看这套源码不只是 Hadoop 和 Spark 那一堆 job,它是一条完整链路:HDFS + Hive 管数据,Spark 管计算,MySQL + Redis 管配置和结果,Spring Boot 管对外服务。任何一个环节拆掉,系统都转不起来。
3. 从 ZIP 到运行:本地复现这套系统的最小启动路径
3.1 环境准备与版本对齐:JDK、Hadoop、Spark、MySQL 一样都不能错
很多人解压完 zip,先看代码,然后mvn package报错,接着就放弃了。跑通这套系统的关键顺序其实是先装环境,再碰代码。我一般会按以下版本组合来做基线,这也是这类源码最保守的选择:JDK 1.8、Hadoop 2.7.6 或 3.2.x、Spark 2.4.x 或 3.1.x、Hive 2.3.x、MySQL 5.7。注意 Spark 2.4 对应 Scala 2.11,Spark 3.x 对应 Scala 2.12,如果源码里用了 Scala 写 UDF,版本不匹配直接编译失败。
# 检查版本对齐,三行命令快速确认 java -version hadoop version spark-shell --version这三个命令的输出版本必须互相匹配。我自己翻过的最离谱的坑是 JDK 装成了 11,Hadoop 3.x 能跑但 Spark 2.4 直接抛UnsupportedClassVersionError,看起来是代码问题,其实纯粹是环境问题。版本对齐这件事没有捷径,建议按源码 README 里给定的版本装,没有 README 的话就按上面这套保守组合来。
3.2 修改核心配置:四份必须改的文件与参数
环境装好后,别急着导入代码,先改配置。Hadoop 端要改core-site.xml、hdfs-site.xml,Spark 端要改spark-defaults.conf,项目自身还有数据库连接配置。这四份文件是跑通的最小集合。
<!-- core-site.xml 里重点确认默认文件系统 --> <property> <name>fs.defaultFS</name> <value>hdfs://localhost:9000</value> </property> <!-- hdfs-site.xml 里伪分布式必须设置副本数为 1 --> <property> <name>dfs.replication</name> <value>1</value> </property>伪分布式模式下dfs.replication必须设成 1,否则集群只有单个 DataNode,副本写入会一直超时。这个参数在源码里经常被带上集群环境的 3,新手装完启动后看着 HDFS 一直Under replicated,其实就是忘了改这里。
Spark 端我一般会关心两个参数:spark.sql.shuffle.partitions和spark.serializer。信贷数据跑 join 时,默认 200 个分区在单机伪分布式里会直接拖垮,改成 8 到 16 更合适;如果源码里用了自定义类在 RDD 之间传递,需要把spark.serializer设为org.apache.spark.serializer.KryoSerializer,否则会报Task not serializable。
数据库配置看项目里的application.yml或jdbc.properties,改成你本地 MySQL 的地址、用户名、密码。这里有个很容易忽略的细节:如果 MySQL 版本是 8.x,JDBC 驱动和连接参数useSSL=false&serverTimezone=Asia/Shanghai不对,后面跑分结果落库时会出现时区错乱,日期差了 13 个小时,回写数据看着像玄学问题。
3.3 导入源码并跑通数据初始化脚本
环境和配置都好了,再导入 IDE。这类工程一般是 Maven 多模块项目,常见有common、etl、risk-engine、api四个模块。导入后先mvn clean install -DskipTests把基础模块装到本地仓库,再逐个编译子模块。编译通过后不要直接调接口,先跑项目里自带的 SQL 初始化脚本。
mysql -uroot -p --default-character-set=utf8mb4 < sql/init.sql执行数据库脚本时我吃过一次亏:在 Linux 终端直接mysql -uroot -p < init.sql,没有指定--default-character-set=utf8mb4,导致规则配置表里的中文描述全部变成乱码。初始化脚本建好库表、插入基础规则配置后,还需要准备一份测试数据放到 HDFS 上,通常是项目里自带的data/目录下的 JSON 或 CSV 样本。
hdfs dfs -mkdir -p /data/ods/apply hdfs dfs -put ./data/apply_20240920.json /data/ods/apply/数据放完,就可以提交前面那段spark-submit命令跑一次规则引擎。判断跑通的标准不是 Spark 作业成功就算了,要看 MySQL 里的风控结果表有没有多出数据、状态字段是不是SUCCESS。只要这一步走出来,整条链路就算打通了。
4. 风控引擎的核心代码怎么改:Spark 读 JSON、规则打分与策略下发
4.1 Spark 读取信贷申请数据:JSON 嵌套结构导致的三类常见问题
信贷申请数据在真实业务里极少是平的。比如一个申请记录里有user_info、loan_apply、credit_report三个嵌套对象,credit_report里还有query_records数组。直接spark.read.json()读进来,DataFrame 会变成带嵌套结构的多层列,你想用select("credit_report.total_loan_amount")没问题,但想去重数组里的某个字段时就会卡住。
from pyspark.sql import SparkSession from pyspark.sql.functions import explode, col spark = SparkSession.builder \ .appName("CreditFeatureJob") \ .enableHiveSupport() \ .getOrCreate() df = spark.read.json("/data/ods/apply/2024-09-20/*.json") df.printSchema() # 展开征信查询记录数组,变成一行一条 query_exp = df.select( col("apply_id"), explode(col("credit_report.query_records")).alias("query_record") ) # 把嵌套字段提取出来做成特征 feature_df = query_exp.select( col("apply_id"), col("query_record.query_date"), col("query_record.query_reason") )这里有一个非常隐蔽的坑:如果 JSON 里某个数组字段在所有记录里都是空的,printSchema()会显示该字段为array<string>但实际上为 null,explode会对 null 数组直接报Can't extract value from explode。源码里如果是老写法,很容易在这一步翻车。解决办法是用explode_outer代替explode,或者先filter(col("query_records").isNotNull())。另外,JSON 里字段名如果有大写或中文,Spark 默认会转成小写并保留反引号,后续所有列名引用都要跟着用反引号包起来,不然后续 join 时列名对不上。
4.2 规则引擎改造:评分卡权重与阈值参数外置到配置表
很多风控系统的规则引擎看起来很高大上,实际上就是若干if-else判断。差的代码把规则写死在 Scala 类里,改一个阈值要重新编译打包;好一点的会把规则抽象成配置项,存在 MySQL 表或 properties 文件里。这套源码我见过比较通用的设计是:规则表存rule_code、feature_code、operator、threshold、score、is_active,规则引擎启动时全量加载到内存,逐条匹配特征值。
// 打分规则加载:从 MySQL 读取规则配置,缓存到本地 Map public class RuleEngine { private Map<String, List<RuleConfig>> rulesByProduct; public void loadRules(DataSource ds) { String sql = "SELECT product_code, rule_code, feature_code, operator, threshold, score, is_active " + "FROM risk_rule_config WHERE is_active = 1"; try (Connection conn = ds.getConnection(); PreparedStatement ps = conn.prepareStatement(sql); ResultSet rs = ps.executeQuery()) { while (rs.next()) { RuleConfig config = new RuleConfig(); config.setFeatureCode(rs.getString("feature_code")); config.setOperator(rs.getString("operator")); config.setThreshold(rs.getBigDecimal("threshold")); config.setScore(rs.getInt("score")); rulesByProduct .computeIfAbsent(rs.getString("product_code"), k -> new ArrayList<>()) .add(config); } } catch (SQLException e) { throw new RuntimeException("加载规则失败", e); } } }这段代码的核心意图是让规则变成数据而不是代码。这样信贷产品经理调额度阈值时,直接执行一条UPDATE risk_rule_config SET threshold = 5000 WHERE rule_code = 'RULE_LOAN_AMOUNT'就行,不需要开发介入。改造时的经验是:外部化规则后,别把is_active做成开关而是做成优先级字段。实际生产里经常有两套规则并行,比如一款产品在灰度期需要新旧策略各跑 50% 流量,如果只支持启停就做不了灰度,必须支持按优先级配置。
4.3 从批处理到结果落库:MySQL 逐条写入必考性能瓶颈
规则引擎跑完 Spark 作业后,结果是一个包含apply_id、final_score、risk_level、hit_rules的 DataFrame。把它写回 MySQL 这一步,新手最容易写出逐条INSERT的代码——跑 50 万条申请时可能要花一个小时。正确做法是用foreachPartition加上 JDBC 批量写入。
import java.sql.{Connection, DriverManager, PreparedStatement} resultDF.foreachPartition { partition => var conn: Connection = null var pstmt: PreparedStatement = null try { conn = DriverManager.getConnection(url, user, pass) val sql = "INSERT INTO risk_result(apply_id, final_score, risk_level, hit_rules, dt) " + "VALUES(?,?,?,?,?) ON DUPLICATE KEY UPDATE final_score=VALUES(final_score)" pstmt = conn.prepareStatement(sql) partition.foreach { row => pstmt.setString(1, row.getAs[String]("apply_id")) pstmt.setInt(2, row.getAs[Int]("final_score")) pstmt.setString(3, row.getAs[String]("risk_level")) pstmt.setString(4, row.getAs[String]("hit_rules")) pstmt.setString(5, bizDate) pstmt.addBatch() if (pstmt.executeBatch().length > 0) pstmt.clearParameters() } } catch { case e: Exception => e.printStackTrace() } finally { if (conn != null) conn.close() } }批量写入的关键是addBatch()累积到一定阈值再executeBatch(),每一条都执行的话性能提升等于零。另一个经验是结果表必须建好唯一键(apply_id+dt),这样重跑任务时可以用ON DUPLICATE KEY UPDATE做幂等更新,而不是先删后插。信贷场景里凌晨批处理失败重跑是常态,没有幂等机制的话,重跑一次就产生一份重复结果,下游审批系统一看同一个申请单出两个分数,直接懵。
5. 部署与调参常见问题排查:5 个踩坑记录与解决路径
5.1 现象:伪分布式切集群模式后作业一直卡在 ACCEPTED
现象:本地伪分布式跑得好好的,换到三台服务器集群后,spark-submit提交的作业一直显示ACCEPTED,既不运行也不报错。
原因:最常见的原因是 YARN 队列资源不够。源码里--executor-memory 8g --num-executors 20是按大集群配的,小集群总内存才 32g,20 个 executor 加上 overhead,每个节点根本塞不下,调度器只能让作业排队等待。
解决:先看yarn-scheduler界面确认可用资源,然后把参数降下来。我一般先按集群物理内存的 60% 估算可分配内存,再除以单 executor 的内存加 overhead 算出 executor 数量。比如 3 台 32g 的机器,单 executor 4g 加 1g overhead,最多跑 11 个 executor,这时候把--num-executors改成 10,--executor-memory改 4g,作业立刻就能跑起来。
5.2 现象:Spark 任务频繁报 OOM,但物理机内存还有大量剩余
现象:跑特征工程聚合时有几个 stage 反复报java.lang.OutOfMemoryError: Java heap space,但去 master 节点看,集群总内存用了不到一半。
原因:这是把 driver 内存和 executor 内存搞混了。某些源码的collect()或take()操作会把大量结果拉回 driver,如果--driver-memory只设了 1g,哪怕 executor 内存有 8g,数据一拉就爆。
解决:先看日志里 OOM 发生在哪一步。如果是collect()报错,把 driver 内存提到 4g,并且检查代码里是否真的需要全量收集,很多时候改成saveAsTextFile或直接落 Hive 更合理。如果是 executor 端 OOM,则要调大spark.executor.memoryOverhead,这个参数是管堆外内存的,默认才executor-memory * 0.1,用 Kryo 序列化时经常不够。
5.3 现象:数据倾斜导致单个 Task 运行特别慢,整批作业卡在一个 stage
现象:跑申请记录和黑名单关联时,整个 Spark 作业其他 task 几秒跑完,就一个 task 跑了几十分钟,最后还会 OOM。
原因:黑名单表里某些身份证号命中次数极高,或者业务分区下某个日期数据量异常大,shuffle 时 hash 落在同一个分区,数据全压在一个 task 上。信贷场景里常见的倾斜源是身份证号 hash、渠道号这种基数低但数据量大。
解决:给 join 的 key 加盐。常见做法是对倾斜 key 做两次 join:先把大表按 key 加随机前缀打散,小表把相同前缀扩展成多份,再进行一次普通 join。
-- 加盐后的大表 SELECT concat(salt, '_', user_id) AS salted_user_id, ... FROM apply_info -- 加盐后的小表 SELECT explode(splits) AS salt, user_id, risk_flag FROM blacklist LATERAL VIEW explode(array('s1','s2','s3','s4')) t AS splits加盐的分桶数要和spark.sql.shuffle.partitions匹配,不然加盐后还是全部冲到同一个 reducer。我一般加 8 到 16 个盐值就够了,加太多反而增加 shuffle 数据量。
5.4 现象:JSON 解析后中文全部乱码,风控规则全变成问号
现象:Spark 读取包含中文渠道名称的 JSON 文件后,查询channel字段显示???,或者规则配置表里的中文从 MySQL 读出来变成乱码。
原因:文件编码和数据库字符集两处问题。JSON 文件可能是 GBK 编码,但 Spark 默认按 UTF-8 读,读出来自然乱码;MySQL 的 JDBC 连接没指定characterEncoding=utf-8,规则加载就乱。
解决:读文件时显式指定编码,连接池参数里加characterEncoding=utf8&useUnicode=true。这里要特别说明,hdfs-site.xml里设置的dfs.replication不会影响编码,编码问题纯粹是读写两端字符集不一致,排查顺序先文件再看库。
# 读文件前先确认编码 file -i /data/ods/apply/2024-09-20/apply.jsonfile -i如果输出charset=utf-8,那问题基本就锁定在 JDBC 连接串和表字段字符集上,顺手查一下SHOW CREATE TABLE risk_rule_config,字段如果还是latin1,直接转utf8mb4。信贷行业风控规则描述里往往会带「禁入」「限制」这种敏感词,如果存储字符集不对,策略管理后台会直接变成一堆问号,看起来很惊悚但其实就是字符集没对齐。
5.5 现象:Hive 表能查到数据但 Spark SQL 查不到同一张表
现象:在 Hive CLI 里能正常SELECT * FROM dwd_apply_info,但 Spark SQL 里跑同一个查询却报Table or view not found。
原因:Spark 默认用的不是 Hive 的元数据库。如果源码里没有配置hive.metastore.uris指向已有的 Hive Metastore,Spark 会在自己的仓库目录里建一套空的元数据,自然看不到 Hive 建的表格。
解决:把 Hive 的hive-site.xml放到 Spark 的conf/目录下,并确保 Spark 启动时能连上hive.metastore.uris指定的端口。还有一个连带的问题是 Hive 元数据库默认存在 Derby 里,多线程并发访问会锁库,生产上一定要把元数据库切到 MySQL。我在本地复现源码时习惯直接把 Hive Metastore 的库放在 MySQL 里,这样调试时不用担心锁库问题。
6. 进阶:从 T+1 批处理到准实时风控的改造路线
批处理跑通只是这个系统的及格线,真实生产环境里很多信贷产品要求小时内响应:新客申请进来,风控结果必须实时出来。改造路线通常是把 Spark 批处理拆两段,一段继续跑天级批量特征,另一段用 Spark Structured Streaming 消费 Kafka 里的申请事件,实时补算轻量特征,再调用规则引擎打分。
关键参数集中在spark.sql.streaming.checkpointLocation和窗口聚合的时间阈值上。批处理里算的是「过去 90 天申请次数」,准实时下不可能真的回溯 90 天,常见做法是预计算好 30 天和 90 天特征存入 Redis,流任务只算当天窗口内的增量特征,两个分数叠加。这里有个时间窗口的取舍,窗口设太短,宽松 5 分钟内的重复申请会漏掉;窗口设太长,实时接口的响应 P99 会直线上升。我一般从 15 分钟起步,根据业务重复申请间隔的样本调参。
验证准实时改造效果不能只看接口延迟,要看一致性:同一批申请数据,用批处理流程跑出的风险和用准实时流程跑出的风险,结果不一致率必须控制在业务容忍范围内。我做过的最直接验证是把最近 7 天的线上申请数据同时喂给两条链路,对比决策结果,不一致率那版差了 3.2%,后来排查发现是流任务读取 Kafka 时默认从最新 offset 开始,丢掉了部分晚到的数据,改成earliest加幂等处理后对齐到 0.3%。
这套系统跳进去容易深挖起来难,后续真正值钱的部分是规则热更新和灰度策略。我做过的教训是:别急着加模型,先把规则配置表和灰度开关做好;很多团队死于能力强但不可控,模型换版本没有灰度,一次线上调参失败能毁掉一整周的客诉指标。按这个顺序往下演进,信我,这套源码能让你在公司里立住脚。希望帮到你。
本文还有配套的精品资源,点击获取