
简介一份基于Hadoop与Spark的大数据金融信贷风险控系统毕业设计源码面向计算机相关专业的在校生、教师或从业者适用于毕设、课设、项目初期立项演示也适合对大数据风控感兴趣的学习者参考进阶。素材包共23个文件压缩后仅53KB以Scala源码为核心配合XML与properties配置文件搭建项目依赖与运行参数另有README说明文档和少量前端静态资源整体轻量紧凑。目前已有2380人学习或浏览积累了一定的参考热度。内容围绕信贷风控场景覆盖Spark流式数据处理组件、数据接入模块及前端展示框架源码经过运行验证功能可用。读者可快速导入开发环境理清HadoopSpark在大数据风控中的应用链路并在此基础上扩展规则引擎或模型算法以完成毕设功能迭代或业务演示从而掌握从数据接入到结果展示的完整流程。无论是用于课程设计还是作为入门大数据工程的项目样例都能提供一个完整可运行的起点帮助读者压缩搭建时间。1. 这份“信贷风控系统”源码包解压之前先想清楚它到底能给你什么很多人看到“HadoopSpark”就觉得这是个大工程其实真正上手之后你会发现毕业设计要的从来不是“造一个银行级风控引擎”而是把一条“数据进得来、算得动、结果出得去”的链路跑通。这个标题里的源代码包大概率是一套这样的骨架HDFS 负责把申请记录和还款流水存下来Hive 做分层Spark 负责清洗、特征加工和风险打分最后落回一张结果表。你能用它答辩、能讲清楚原理、能改参数演示这就够了。如果你是准备大数据方向毕业设计、或者想转大数据开发但手里缺一个完整项目的人这套方向是值得花时间吃透的。但别指望解压就能跑环境、版本、内存这些坑都在后面等着。2. 信贷风控为什么要押注 Hadoop Spark先把数据链路画清楚2.1 金融信贷风控的离线批处理链路从申请记录到风险分信贷风控在真实业务里分贷前、贷中、贷后毕业设计一般只做贷前这一环用户提交贷款申请系统根据年龄、收入、负债、历史逾期次数、申请金额等字段给一个风险评分再决定放不放款。这个场景特别适合用大数据技术栈不是因为数据量真的到了“海量”而是因为特征计算逻辑复杂、数据来源多、需要批量调度用 HadoopSpark 能把“多份数据 → 一张评分表”的加工过程做得规规矩矩。常见的数据链路是四层先是原始数据层存申请记录、流水、第三方征信数据对应 ODS 层再是清洗层把缺字段、异常值、格式错乱的记录处理掉对应 DWD 层然后是汇总层按客户维度聚合出特征对应 DWS 层最后是应用层输出每个客户的风险分和风险等级对应 ADS 层。你拿到任何一份这样的毕设源码都可以按这四层去找它的代码结构找不到就说明这个毕设本身分层不清楚。Hadoop 在这条链路里的角色是“底座”HDFS 管存储Hive 管数仓MapReduce 虽然也能算但写起来太痛苦所以重活都交给 Spark。Spark 拿手的正好是内存计算和 DAG 调度无论是读 JSON 原始文件做 ETL还是做客户维度聚合、特征宽表拼接都比 MapReduce 快一个数量级代码也短得多。对这个项目而言Spark 就是那个“算”的角色Hadoop 是那个“存”的角色。2.2 Hadoop 在风控里的角色HDFS 目录规划与 Hive 分层HDFS 上怎么规划目录决定了你后续所有脚本好不好写。常见的做法是建一个统一根目录比如/risk下面分/risk/ods、/risk/dwd、/risk/dws、/risk/ads每层下面再按日期分区。这样做的好处是跑批任务时只要指定“读哪一层、写哪一层”不会出现数据乱窜。对应的 Hive 建表语句大概长这样这张表是用来放 ODS 层原始申请记录的CREATE EXTERNAL TABLE ods_credit_apply ( customer_id STRING, apply_date STRING, age INT, annual_income DOUBLE, debt DOUBLE, overdue_times INT, apply_amount DOUBLE ) PARTITIONED BY (dt STRING) ROW FORMAT DELIMITED FIELDS TERMINATED BY , STORED AS TEXTFILE LOCATION hdfs://localhost:9000/risk/ods/credit_apply;用 EXTERNAL TABLE 而不是内部表是因为原始数据经常是外部脚本生成的删表不能把数据文件也删了。PARTITIONED BY (dt STRING)这个分区字段很关键你后面跑 Spark 任务时按天增量处理就是靠这个 dt 来限定数据范围。LOCATION指向的路径要和第 3 章里 HDFS 的fs.defaultFS对应上。在真实项目里ODS 层通常会比上面这张表多很多字段渠道来源、设备指纹、IP 归属地、填表耗时等等这些字段对反欺诈有用。毕业设计不追求这么多但你要在答辩时说出“为什么 ODS 层要尽可能保留原始字段、不做过滤”——因为过滤逻辑在下游做原始数据一旦丢了就补不回来。这句话能体现你理解数仓分层的用意。2.3 Spark 在风控里的角色ETL 清洗、特征加工、风险打分三件事Spark 在这个系统里不是只跑一个任务而是按职责拆成三条线第一条是 ETL 清洗把 ODS 层 JSON 转成结构化 DataFrame处理空值和异常值第二条是特征加工把清洗后的明细按客户维度做聚合比如算“过去 6 个月平均申请金额”“负债收入比”“历史逾期次数分段”这类特征第三条是风险打分把特征加权求和或者套一个简单的评分卡得到 risk_score 和 risk_level。这三个阶段在代码上最好拆成三个独立的方法或者三个独立的 Spark Job不要揉在一起。原因很实际你调参的时候比如想改打分权重只需要重新跑第三段前面清洗结果还在 Hive 表里存着不用从头再来。如果揉成一个 Job每次改一个系数就要把全链路重跑一遍程序是能跑但你会浪费大量等待时间。也要想清楚什么时候用 SparkSession、什么时候用 SparkContext。现在统一用 SparkSession 就行它把 SQL、DataFrame、Streaming 的入口都收拢了。跑批任务里要注意的是 driver 和 executor 的分工driver 负责调度和收集结果executor 负责实际计算。如果你把几百万条结果都 collect 回 driver内存直接爆掉。这也是为什么写完评分以后正确姿势是write.saveAsTable落回 Hive而不是打印到控制台。3. 把 Hadoop Spark 跑起来伪分布式搭建与最小集群参数3.1 环境选择伪分布式够用吗先算清楚你的数据量很多人在环境搭建这一步就卡了两三天原因不是技术难而是没想清楚“我到底需要多大环境”。毕业设计的数据量撑死几百万条单机伪分布式完全扛得住。所谓伪分布式就是一台机器上同时跑 NameNode、DataNode、ResourceManager、NodeManager 这四个进程HDFS 的副本数设为 1。它和真正集群的区别只在“副本数”和“进程分布”API 用法一模一样。如果你只想在 Windows 笔记本上把流程跑通可以装虚拟机或者直接用 Docker 起 Hadoop 镜像但如果你后面要写“三节点集群搭建”这种章节来凑字数那就老老实实准备三台虚拟机。我的建议是不要把时间浪费在搭集群上伪分布式能让你把 80% 精力放在业务代码和答辩逻辑上。面试官问“集群部署策略”你能说出“生产环境 NameNode 和 ResourceManager 要分离部署、DataNode 按机架划分”就足够了毕设演示机器上不需要真做。版本搭配也是一个大坑。常见问题是拿着 Hadoop 2.x 的教程配 Hadoop 3.x或者 Spark 版本和 Hadoop 版本编译不兼容。毕业设计我一般建议这套组合Hadoop 3.2.x Spark 3.xJDK 用 8 或 11。Spark 3.x 对 Hive 的支持更顺滑saveAsTable直接写 Hive 表不容易踩序列化坑。不要追新版本JDK 17 最新 Hadoop 的组合在配置上会有额外麻烦。3.2 伪分布式最小配置core-site.xml、hdfs-site.xml、yarn-site.xml 三件套Hadoop 跑起来只需要改三个配置文件都在$HADOOP_HOME/etc/hadoop目录下。第一个是core-site.xml指定默认文件系统和临时目录configuration property namefs.defaultFS/name valuehdfs://localhost:9000/value /property property namehadoop.tmp.dir/name value/opt/hadoop/tmp/value /property /configurationfs.defaultFS是所有 HDFS 路径的前缀后面代码里写hdfs://localhost:9000/risk/input/apply.json就是在用它。hadoop.tmp.dir是 NameNode 和 DataNode 存放元数据与数据块的根目录如果这个目录没有创建权限或者磁盘满了后面格式化会失败。第二个是hdfs-site.xml核心是把副本数设为 1并显式指定 NameNode 和 DataNode 的数据目录configuration property namedfs.namenode.name.dir/name valuefile:///opt/hadoop/namenode/value /property property namedfs.datanode.data.dir/name valuefile:///opt/hadoop/datanode/value /property property namedfs.replication/name value1/value /property /configuration这里有一个很容易被忽略的点dfs.namenode.name.dir和dfs.datanode.data.dir如果配置的父目录不存在Hadoop 不会自动创建启动会直接报错。你需要在配置文件改好以后手动mkdir -p /opt/hadoop/namenode和mkdir -p /opt/hadoop/datanode。第三个是yarn-site.xml资源调度全靠它。伪分布式下最容易翻车的是内存配置过大笔记本根本给不起configuration property nameyarn.nodemanager.resource.memory-mb/name value8192/value /property property nameyarn.scheduler.maximum-allocation-mb/name value4096/value /property property nameyarn.nodemanager.resource.cpu-vcores/name value4/value /property /configurationyarn.nodemanager.resource.memory-mb是 NodeManager 能使用的总内存yarn.scheduler.maximum-allocation-mb是单个容器能申请的最大内存。如果你的机器只有 8G 内存建议把前者改成 4096后者改成 2048否则启动以后 Spark 任务很容易因为申请不到资源而一直卡在 ACCEPTED 状态。3.3 一键拉起与验证格式化、启动、看进程、跑第一个 Spark SQL配置文件改完以后的启动顺序是有讲究的先格式化 NameNode再启动 HDFS再启动 YARN。格式化只要一次第二次启动不用再格式化这是新手最容易犯的错——每次重新格式化会导致元数据丢失DataNode 和 NameNode 的集群 ID 对不上启动直接失败。# 第一次使用必须格式化 NameNode hdfs namenode -format # 启动 HDFS start-dfs.sh # 启动 YARN start-yarn.sh # 确认四个进程都在 jpsjps输出里应该看到四个名字NameNode、DataNode、ResourceManager、NodeManager。少一个都说明配置有问题这时候去查对应日志重点看$HADOOP_HOME/logs目录下的hadoop-hadoop-namenode-*.log和hadoop-hadoop-datanode-*.log。启动完成以后别急着跑业务代码先验证 Spark 能和 HDFS、YARN 正常通信。最简单的办法是提交一个 Spark SQL 任务让它读 HDFS 上的一个小文件# 先往 HDFS 放一个测试文件 echo hello hadoop | hdfs dfs -put - /test.txt # 用 spark-sql 读 HDFS 文件能返回结果就说明链路通了 spark-sql --master yarn --deploy-mode client \ -e SELECT count(*) FROM (SELECT 1 AS id UNION ALL SELECT 2) t;这里用--master yarn让 Spark 跑在 YARN 上而不是本地模式。很多人的毕设代码在本地能跑一到集群就挂就是因为 master 没配成 yarn。等这一步通了就可以进入下一章开始处理真正的信贷数据。4. 让源码包里的风控逻辑跑通业务闭环数据造数、Spark 读取 JSON、打分落库4.1 造一套信贷模拟数据Python 写申请记录生成器如果你拿到的源码包没有配套的数据文件那就必须自己造一套。不要用手工敲的几十条数据那个量级跑 Spark 没有说服力。我一般用 Python 写一个生成器把客户 ID、年龄、年收入、负债、历史逾期次数、申请金额写成一行为一条 JSON 的文件模拟真实申请记录。import json import random from datetime import datetime, timedelta random.seed(42) rows [] for i in range(100000): customer_id fC{100000 i} age random.randint(20, 60) annual_income random.randint(60000, 1200000) debt random.randint(0, int(annual_income * 0.8)) # 大部分客户没有逾期少数逾期 1-3 次 overdue_times random.choices([0, 1, 2, 3], weights[0.80, 0.12, 0.05, 0.03])[0] apply_amount random.randint(10000, 500000) rows.append({ customer_id: customer_id, apply_date: (datetime(2024, 1, 1) timedelta(daysrandom.randint(0, 365))).strftime(%Y-%m-%d), age: age, annual_income: annual_income, debt: debt, overdue_times: overdue_times, apply_amount: apply_amount }) with open(credit_apply.json, w, encodingutf-8) as f: for record in rows: f.write(json.dumps(record, ensure_asciiFalse) \n) print(fgenerated {len(rows)} rows)random.seed(42)保证每次生成的数据一样这一点在答辩时很重要——评审老师如果让你现场重跑你不想让结果每次都变。overdue_times用random.choices按权重抽样让 80% 客户是 0 次逾期、少数客户 1 到 3 次这样后面打分评出来的风险等级才符合业务直觉大部分是低风险小部分是高风险。生成的是 JSON Lines 格式一行一条记录Spark 可以直接读。4.2 Spark 读取 JSON 做特征清洗用 DataFrame API 的完整代码数据生成以后把它传到 HDFS 上然后编写 Spark 清洗任务。这一步的目标是读入 JSON → 过滤_corrupt_record→ 剔除年龄、收入不合理的记录 → 计算年月分区字段 → 计算负债收入比特征 → 写入 Hive 表。import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ val spark SparkSession.builder() .appName(credit-risk-etl) .enableHiveSupport() .getOrCreate() val raw spark.read .option(mode, PERMISSIVE) .option(columnNameOfCorruptRecord, _corrupt_record) .json(hdfs://localhost:9000/risk/input/credit_apply.json) val cleaned raw .filter(col(_corrupt_record).isNull) .drop(_corrupt_record) .filter(col(age).between(18, 65)) .filter(col(annual_income) 0) .withColumn(apply_year_month, date_format(to_date(col(apply_date)), yyyyMM)) .na.fill(0, Seq(overdue_times)) .withColumn(debt_ratio, round(col(debt) / col(annual_income), 4)) cleaned.write .mode(overwrite) .partitionBy(apply_year_month) .format(parquet) .saveAsTable(dwd_credit_apply)这段代码的核心是mode(PERMISSIVE)它会把解析失败的记录塞进_corrupt_record字段而不是让整个任务直接崩掉。紧接着.filter(col(_corrupt_record).isNull)把脏数据过滤掉再.drop删掉这个辅助字段。这一步如果你做反了——先 drop 再 filter——脏数据就被丢掉了但你也永远不知道到底有多少脏数据这对风控场景来说是致命的。.partitionBy(apply_year_month)是让结果在 Hive 里按月份分区存储后续查“某个月有多少高风险客户”时只需扫描一个分区。saveAsTable会同时建表和写数据如果表已存在mode(overwrite)会覆盖数据但注意它默认不会替换表结构。4.3 规则引擎打分与结果写回 Hive从 risk_score 到可用报表清洗后的表在dwd_credit_apply下一步是打分。毕业设计里的模型一般不会用真正训练出来的机器学习模型而是用简化的评分卡把收入、逾期次数、负债收入比三个特征分别映射成分值加总得到risk_score再按阈值切成低中高三个风险等级。val scored cleaned .withColumn(income_score, when(col(annual_income) 500000, 30) .when(col(annual_income) 200000, 20) .otherwise(10)) .withColumn(overdue_score, when(col(overdue_times) 0, 40) .when(col(overdue_times) 1, 20) .otherwise(0)) .withColumn(debt_score, when(col(debt_ratio) 0.3, 30) .when(col(debt_ratio) 0.6, 15) .otherwise(0)) .withColumn(risk_score, col(income_score) col(overdue_score) col(debt_score)) .withColumn(risk_level, when(col(risk_score) 80, low) .when(col(risk_score) 50, medium) .otherwise(high)) scored.write .mode(overwrite) .format(parquet) .saveAsTable(ads_credit_risk)这里的三个when就是一张最简单的评分卡权重完全是可调的。你把收入阈值从 50 万调到 30 万高分客户的数量立刻变多。答辩时你可以现场演示“调高收入权重之后高风险客群占比从 15% 降到了 9%”这是很加分的操作。打分完成以后任务往往需要提交到集群上跑。在源码包里通常是一个打好的 jar你只需要把主类和输入输出路径换成自己的spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 2g \ --num-executors 2 \ --executor-cores 1 \ --driver-memory 1g \ --class com.risk.CreditRiskApp \ credit-risk-1.0.jar \ hdfs://localhost:9000/risk/input/credit_apply.json \ hdfs://localhost:9000/risk/result--executor-memory 2g是指每个 executor 分配 2G 堆内存--num-executors 2是启动两个 executor。这两个参数在伪分布式下不要贪大前面 yarn-site.xml 里的maximum-allocation-mb配置是多少executor 内存就不能超过多少否则任务会一直卡在 ACCEPTED 状态。--deploy-mode cluster会让 driver 跑在集群里关掉客户端也不影响任务执行适合放在nohup后面挂后台。5. 实战踩坑HadoopSpark 跑信贷风控任务的 5 个高频故障5.1 NameNode 起不来jps 里永远少一个进程现象执行start-dfs.sh后用jps只能看到 DataNode 和 ResourceManagerNameNode 进程不存在查看日志发现报错信息里有file:///opt/hadoop/namenode has been used by another process或者Incompatible clusterIDs。原因最常见的是两种一是/opt/hadoop/namenode目录里已经有了旧元数据你第二次执行了hdfs namenode -format生成了新的 cluster ID和 DataNode 里记录的 cluster ID 不一致二是端口 9000 被占用NameNode 启动时绑定端口失败。解决如果没有重要数据最简单的方法是彻底清理之后重新格式化一次。把hadoop.tmp.dir下所有内容删除包括 namenode、datanode、tmp 三个目录然后重新mkdir -p建目录再执行hdfs namenode -format。千万不要保留旧目录直接格式化。如果是端口占用用ss -lntp | grep 9000找到占用进程处理掉再重启。5.2 Spark 任务提交后一直卡在 Accepted 状态等十分钟都不跑现象spark-submit执行后客户端一直打印INFO YarnClientImpl: Submitted application ...然后状态停在ACCEPTED没有 Executor 启动的日志。原因申请不到资源。要么是yarn.nodemanager.resource.memory-mb小容器申请的内存大于上限要么是--executor-memory加--driver-memory的总和超过了 NodeManager 可用内存YARN 把 ApplicationMaster 都调度不出来。解决先去 YARN 的 Web UIhttp://localhost:8088看集群内存总量和可用量确认没看错以后把--executor-memory改小到 1g--num-executors改成 1再把 yarn-site.xml 里的yarn.scheduler.maximum-allocation-mb调到不低于 2g。伪分布式环境下资源和代码一样重要调参本来就是毕设的一部分。5.3 按客户聚合特征时数据倾斜某个 reduce 直接 OOM现象跑“按 customer_id 聚合申请金额、次数”这一类groupBy操作时大部分 task 秒完成但有个别 task 要跑几分钟甚至报GC overhead limit exceeded。看 Spark UI 的 Stage 页某个 Executor 的 Shuffle Read 量是其他 Executor 的几十倍。原因信贷数据里存在“羊毛党”或“大客户”一个人可能提交了几万次申请生成数据时如果没控制分布customer_id这一个 key 就占了大头。这个 key 所在的 reducer 被迫处理远高于平均的数据量。解决做两阶段聚合先加盐拆开热点 key聚一次再去掉盐做第二次聚合。代码很短val salted detail .withColumn(salt, (rand() * 10).cast(int)) .withColumn(group_key, concat(col(customer_id), lit(_), col(salt))) val partial salted .groupBy(group_key) .agg(sum(apply_amount).as(partial_amount)) val finalResult partial .withColumn(customer_id, split(col(group_key), _).getItem(0)) .groupBy(customer_id) .agg(sum(partial_amount).as(total_amount))加盐字段只影响中间分组最终结果不受影响。盐的粒度可以调10 不够就加到 100但不要无限加否则小文件也会变多。这个方案不只能用在风控场景任何热点 key 导致的倾斜都能用。5.4 JSON 解析失败导致整个批次静默丢数据现象清洗任务跑完dwd_credit_apply表里记录数比 ODS 层少了上千条但日志里没有任何 ERROR连 warning 都很少。仔细看才发现某些钱多的客户被过滤掉了。原因spark.read.json默认的mode是PERMISSIVE没错但如果你没有设置columnNameOfCorruptRecord并且在后续代码里显式 filter脏数据会被直接 discard而且不报错。你的造数脚本里如果出现了中文字段值或者空数组就可能触发这个问题。解决用上一章 4.2 里的写法设置columnNameOfCorruptRecord为_corrupt_record然后必须紧跟在 read 之后马上 filter再用negate或isNull判断。事后补救的办法也有如果数据已经写进 Hive 表可以用spark.read.table(dwd_credit_apply).filter(_corrupt_record is null)再把脏数据筛掉但更靠谱的做法是改代码重跑。5.5 写入 Hive 分区表后小文件爆炸查询慢到怀疑人生现象任务正常结束但dwd_credit_apply表按apply_year_month分区每个分区下有几百个几 KB 大小的文件SELECT count(*)要跑几十秒。原因Spark 写文件时每个 shuffle 分区都会为每个目标分区写文件。如果spark.sql.shuffle.partitions默认是 200而你的数据只有几十万行每个文件都特别小。Hive 读大量小文件比读一个大文件慢得多。解决写完以后合并小文件或者在写之前控制分区数。最简单的是重跑时加一行配置spark-sql --master yarn \ --conf spark.sql.shuffle.partitions10 \ -e INSERT OVERWRITE TABLE dwd_credit_apply PARTITION(apply_year_month) SELECT ... FROM ods_credit_applyspark.sql.shuffle.partitions设置为 10能让每个分区只写一个文件。已经产生的小文件可以用ALTER TABLE dwd_credit_apply PARTITION(apply_year_month202401) CONCATENATE;在 Hive 里合并掉。你可以在答辩时说“如果不控制 shuffle 分区数小文件会让 NameNode 内存吃紧”——这句话是加分项说明你懂小文件对 Hadoop 的危害。6. 把毕设从“能跑”变成“能讲”两个改造方向和一个验证技巧6.1 改造方向一给 HDFS 加 ZooKeeper 高可用让系统更完整伪分布式做的是单 NameNode面试时被问到“NameNode 挂了怎么办”你要能接住这个话题。常见的做法是引入 ZooKeeper 做 NameNode HA两个 NameNode 一主一备共享编辑日志存在 JournalNode 里ZooKeeper 负责故障切换。这个改造不用在毕设演示机上真的做但你要会画架构图、能说清楚 ZK 选主的过程。给源码加一个ha分支或配套文档模板比硬写代码更实际。这里踩过的坑是ZooKeeper 的 myid 写错会导致选主失败日志里全是cant connect to quorum排错时先zkServer.sh status看每个节点的 state再检查防火墙和地址配置。6.2 改造方向二用 Spark Streaming 做一个风险预警模块纯离线的信贷风控系统只能做贷前审批如果能加一个实时模块比如监控还款流水当某个客户出现“短时间多笔大额取现”时实时拉高风险等级整个项目的完整度会上一个台阶。Spark Streaming 现在的写法是用 Structured Streaming读 Kafka 或文件源窗口聚合判断异常写回 Hive 或 Redis。在毕设里你不需要接真数据用文件源加上一个模拟数据脚本就能演示。6.3 验证技巧用影子对照实验证明评分有效答辩时最怕的是被问“你怎么证明你的评分是有效的”。不需要上机器学习那套 AUC你只要做一个简单的分布验证分别计算 low、medium、high 三组客户的平均逾期次数然后对比。spark-sql -e SELECT risk_level, count(*) AS cnt, round(avg(overdue_times), 2) AS avg_overdue FROM dwd_credit_apply GROUP BY risk_level ORDER BY risk_level; 如果 high 组的avg_overdue明显高于 low 组说明你定的阈值有区分度。反过来说如果三组的平均逾期次数几乎一样说明你的评分卡权重是拍脑袋拍的回去调权重再跑一次。我当年做这个项目时最吃亏的一件事就是只顾着把链路跑通没准备“证明它有效”这一步答辩被问住以后才补的这套对照实验。这个习惯我现在做任何数据项目都会保留——先想好怎么证明结果是对的再动手写代码。希望帮到你。本文还有配套的精品资源点击获取