☰
基于Hadoop和Spark的信贷风控大数据离线链路构建
2026/10/2 3:08:50 网站建设 项目流程

简介:基于Hadoop和Spark的金融信贷风控大数据系统毕业设计源码,主要面向计算机专业学生与大数据实践者,适用于毕业设计、课程作业或项目实战,重点解决海量信贷数据下的风险预测与实时监控问题。压缩包共69个文件,大小约69KB,核心代码由36个Java和8个Scala源文件构成,另有XML、properties、SQL脚本及Markdown文档辅助,极大方便快速了解工程配置与数据表结构设计。已有76人学习下载。项目经过导师指导并完成本地编译调试,源码可正常运行,且经过严格调试,整体覆盖数据采集、预处理、风险评估模型构建与结果输出等完整流程;同时按主项目与数据接入模块分块组织,附带数据库脚本和说明文档,非常适合参考系统架构、二次开发或作为大数据金融风控实战练习,能有效缩短从理论到落地的学习成本。

1. 信贷风控走到大数据这一步:这个毕设课题到底在做什么

把一笔贷款放出去能不能收回来,信贷风控靠的是对用户行为的判断;一天进来几十万笔借款申请,靠的则是一套基于 Hadoop 和 Spark 的信贷风控大数据系统。这个题目在毕业设计里属于标准的大数据离线链路题:从 HDFS 上的原始流水出发,用 Spark SQL 清洗加工,聚合出用户维度的风险特征,再交给模型产出风险分。真正卡住大多数人的不是模型算法,而是怎么把这条完整的数据链路跑通、讲清楚、并且扛得住答辩追问。这篇笔记写给两类人:一类是正在做大数据毕业设计、为选题和代码发愁的学生;另一类是刚接触离线数仓、想用一个完整案例练手 Spark 的开发。读完你能拿到一套可以直接改表名的代码骨架,以及环境搭建、参数设置和验证阶段常用的避坑方法。

2. 先选型再动手:信贷风控大数据的分层架构与 Hadoop/Spark 分工

信贷风控系统最忌讳把所有计算逻辑塞进一个程序里一把梭。业界做这类离线系统的常规思路是按数据流向分层,而不是按业务功能切模块。分层的价值在于每一层只做一件事,出了问题能定位到具体环节,重跑成本也可控。整个链路下来,你写出来的每个代码文件都应该能说出它属于哪一层。

2.1 数据分几层:贴源层、清洗层、加工层、应用层分别解决什么

一张典型的信贷风控大数据离线链路可以切成四层:

层常用命名职责落地组件
贴源层ODS原始贷款申请、放款流水、还款流水原样落盘,出错可重放HDFS 目录 + Parquet
清洗层DWD去重、格式统一、异常值过滤,生成明细事实表Spark SQL
加工层DWS按用户/时间维度聚合,生成特征宽表Spark SQL + DataFrame
应用层ADS风险分、统计看板、预警名单Spark MLlib + MySQL

毕设里最容易拿分的是 DWS 层,因为特征宽表直接决定后续模型效果。很多同学把清洗和特征混在一个脚本里,论文阶段会很痛苦,因为答不出“这条数据为什么在这里被过滤掉”。如果答辩老师问“你的系统分几层”,你至少要把上面这张表讲清楚:每层输入什么、输出什么、用什么组件跑的。

对应到目录,我建议你在 HDFS 上建四个根目录:/warehouse/ods、/warehouse/dwd、/warehouse/dws、/warehouse/ads。代码里读写路径统一从/warehouse开始,后期找数据、增量重跑都方便。

2.2 为什么是 Hadoop 存、Spark 算:HDFS、YARN、Spark 引擎的取舍

信贷流水的数据量一上去,单机数据库先撑不住的不是计算,是存储和 IO。HDFS 的价值在于把几十 GB 甚至 TB 级文件分散到多台机器上,靠副本策略保证数据不丢,不需要做 RAID。而 YARN 负责资源调度,Spark 作业向 YARN 申请容器跑 executor。风控特征加工的核心操作是 groupBy 和 join,这类任务对 shuffle 的依赖很强,Spark 把中间结果留在内存里,比 MapReduce 一轮轮落盘快得多,这也是为什么选 Spark 而不是传统 MapReduce。

实际生产里,这套组合最常见的形态是:数据落 HDFS,用 Spark SQL 做离线批处理,结果写 Hive 表或 MySQL,供风控后台查询。很多公司还会在同一个 YARN 集群上跑 Flink 实时任务,但那是另一条链路了,毕业论文里不用展开。如果数据量只有几百万行,用 MySQL 加 Pandas 也能跑,没必要上 Hadoop;反过来,如果选题就是“大数据风控”,你不把 HDFS、YARN、Spark 这三个组件同时用上,答辩时容易被一句话问倒:你的大数据体现在哪里。

2.3 落地环境怎么搭:伪分布式 Hadoop 与 Spark on YARN 的关键配置

环境选型只有两种常见路线。一是单机伪分布式,适合时间紧、机器配置一般的同学;二是三台机器的小集群,适合想展示“集群部署能力”的同学。伪分布式不是“假”的,它让 NameNode、DataNode、ResourceManager 都跑在同一台机器上,进程是齐的,代码不用改,只是并发能力弱。以下配置项直接写进core-site.xml、hdfs-site.xml、yarn-site.xml和spark-defaults.conf:

配置文件配置项建议值作用
core-site.xmlfs.defaultFShdfs://localhost:9000指定 NameNode 地址
hdfs-site.xmldfs.replication1(伪分布式)/ 3(集群)副本数
yarn-site.xmlyarn.nodemanager.resource.memory-mb机器内存的 60%~70%NodeManager 可用内存上限
yarn-site.xmlyarn.scheduler.maximum-allocation-mb与上面接近单个容器最大内存
spark-defaults.confspark.masteryarn让 Spark 任务跑在 YARN 上
spark-defaults.confspark.executor.memory1g~2g每个 executor 堆内存

下载 Hadoop 和 Spark 的二进制包后,按下面顺序操作,先启动 HDFS,再启动 YARN,然后提交一个测试任务确认链路通:

# 解压到 /opt 后,先配置 JAVA_HOME 和 PATH,再格式化 NameNode export JAVA_HOME=/opt/jdk export PATH=$PATH:$JAVA_HOME/bin:/opt/hadoop/sbin:/opt/hadoop/bin:/opt/spark/bin # 首次使用必须格式化,之后不要再执行,否则丢元数据 hdfs namenode -format # 启动 HDFS 和 YARN sbin/start-dfs.sh sbin/start-yarn.sh # 确认进程:NameNode、DataNode、ResourceManager、NodeManager 都在 jps

启动完成后用 jps 检查进程,漏了哪个就去对应日志看报错。之后跑一个 Spark 官方自带的计算 Pi 任务验证 YARN 调度正常:

spark-submit --master yarn --deploy-mode client \ /opt/spark/examples/jars/spark-examples_*.jar 100

deploy-mode 用 client 方便看日志,伪分布式环境不需要 cluster 模式。如果任务能正常算出结果,说明 Spark 和 YARN 已经通了,后面写的代码只需要按同一套 submit 参数提交。

3. 用 Spark SQL 把流水变成特征:信贷风控宽表加工与逾期标签

这一章是整个项目的核心工作量所在。信贷风控的特征加工说白了就是回答几个问题:这个人最近借了多少次、借了多少钱、有没有逾期、当前还欠着多少。所有特征最终落成一张以 user_id 为主键的宽表,供下游模型使用。为了不引入额外的 Hive 服务,这里直接用 Spark SQL 读写 Parquet 文件,Parquet 文件目录本身就是表。

3.1 先造一份最小信贷数据集:表结构与模拟数据

没有真实数据时,写一个造数脚本生成三张表:用户表、借款申请表、还款流水表。字段不要贪多,够建模和答辩演示就好。

数据文件关键字段
users.csvuser_id, name, gender, reg_time
loans.csvloan_id, user_id, apply_time, amount, term, status
repays.csvrepay_id, loan_id, user_id, repay_time, amount, overdue_days

直接用 Python 生成 CSV,再让 Spark 读入转成 Parquet 落盘:

import csv import random import datetime random.seed(42) user_count = 20000 loan_count = 50000 # 生成用户表:2 万用户,注册时间分布在近三年 with open("users.csv", "w", newline="") as f: writer = csv.writer(f) writer.writerow(["user_id", "name", "gender", "reg_time"]) for i in range(1, user_count + 1): reg = datetime.date(2018, 1, 1) + datetime.timedelta(days=random.randint(0, 1000)) writer.writerow([i, f"user_{i}", random.choice(["M", "F"]), reg]) # 生成借款流水:每笔借款有金额、期限、状态 with open("loans.csv", "w", newline="") as f: writer = csv.writer(f) writer.writerow(["loan_id", "user_id", "apply_time", "amount", "term", "status"]) for i in range(1, loan_count + 1): apply = datetime.date(2019, 1, 1) + datetime.timedelta(days=random.randint(0, 700)) writer.writerow([ i, random.randint(1, user_count), apply, random.randint(2000, 200000), random.choice([3, 6, 12]), random.choice(["已结清", "逾期", "还款中"]) ])

这段代码里random.seed(42)保证每次生成的数据一致,方便答辩时复现结果。用户量不用大,2 万用户配 5 万笔借款已经是百万行以下的小数据,跑起来快,后续讲 Spark 的优势再准备一个放大版数据就行。

Spark 读取并转成 Parquet 的代码如下,落盘后/warehouse/dws下面就是你自己的“表”:

from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("init_data") \ .master("yarn") \ .getOrCreate() df_users = spark.read.option("header", True).csv("users.csv") df_loans = spark.read.option("header", True).csv("loans.csv") df_users.write.mode("overwrite").parquet("/warehouse/ods/users") df_loans.write.mode("overwrite").parquet("/warehouse/ods/loans")

mode("overwrite")保证重复跑脚本不会报“目录已存在”的错误。Parquet 是列式存储,压缩率高,Spark 读起来比 CSV 快很多,生成一遍以后全程用 Parquet。

3.2 数据清洗:过滤异常流水与重复申请

原始流水里最常见的三个问题:user_id 为空、借款金额明显异常、同一用户同一天重复申请。清洗的目标是让明细层至少逻辑自洽。

from pyspark.sql.functions import col, to_date df_loans_clean = ( spark.read.parquet("/warehouse/ods/loans") .filter(col("user_id").isNotNull()) .filter(col("amount").between(100, 1000000)) .withColumn("apply_date", to_date(col("apply_time"))) .dropDuplicates(["user_id", "apply_date", "amount"]) .write.mode("overwrite").parquet("/warehouse/dwd/loans_clean") )

四个步骤各有用途:isNotNull 过滤掉脏数据;between 把金额在 100 元以下或 100 万以上的申请视为异常,这两种极端值后续会拉偏特征;to_date 把字符串时间转成日期类型,后面算“近 30 天”这类时间窗口时可以直接比较;dropDuplicates 按 user_id、apply_date、amount 去重。

注意 dropDuplicates 默认保留第一行,不代表保留最新的状态。如果需要保留最新状态,先按 apply_time 降序排序再 dropDuplicates。大多数场景下,重复申请本身就是一个风险特征,一小时内连续申请多次的用户更需要被标记出来,所以这一步只要去重,不需要额外删人。

3.3 特征聚合:从借贷流水到用户画像宽表

宽表加工的核心是按 user_id 分组,把多笔借款流水压成一个用户的多列特征。这里给出两个最常用的聚合口径:近 30 天申请强度、历史逾期情况。

from pyspark.sql.functions import count, sum, max, when, datediff, current_date df_loans = spark.read.parquet("/warehouse/dwd/loans_clean") # 近 30 天申请次数与金额:注意先过滤时间窗口,再做 groupBy df_recent30 = ( df_loans .filter(col("apply_date") >= datediff(current_date(), 30)) .groupBy("user_id") .agg( count("loan_id").alias("recent_30_apply_cnt"), sum("amount").alias("recent_30_amount") ) ) # 历史逾期与在贷特征:不需要时间窗口,全量累计 df_overdue = ( df_loans .groupBy("user_id") .agg( count(when(col("status") == "逾期", 1)).alias("overdue_cnt"), max(col("amount")).alias("max_loan_amount"), sum(when(col("status") == "还款中", col("amount"))).alias("current_loan_amount") ) )

把两张聚合表 join 成宽表时,必须以用户表为左表,用 left join,防止没有借贷记录的用户被丢掉:

df_wide = ( spark.read.parquet("/warehouse/ods/users") .join(df_recent30, "user_id", "left") .join(df_overdue, "user_id", "left") .fillna(0) ) df_wide.write.mode("overwrite").parquet("/warehouse/dws/user_features")

这段代码里的fillna(0)很关键。left join 之后没有记录的用户会出现空值,模型训练时空值会直接报错。用 0 填充符合业务理解:没借过钱的人,申请次数和逾期次数就是 0。特征完整后,宽表就是最终模型的训练输入。

3.4 必调参数:分区数、内存与动态裁剪

Spark 作业跑得慢,十有八九不是代码逻辑问题,是参数没跟上。先记住三个最常用的:

参数默认值这个项目建议说明
spark.sql.shuffle.partitions20050伪分布式或小数据量,200 个分区反而浪费调度开销
spark.executor.memory1g1g~2g超过 NodeManager 可用内存会卡在 ACCEPTED
spark.sql.adaptive.enabledfalsetrueSpark 3 开启 AQE,自动合并小分区

数据只有几万行时,shuffle 分区改成 50 甚至 20 能让作业快很多。提交时这样指定:

spark-submit --master yarn \ --conf spark.sql.shuffle.partitions=50 \ --conf spark.sql.adaptive.enabled=true \ --conf spark.executor.memory=2g \ loan_feature.py

如果之后把数据放大到百万行,shuffle 分区再回调到 100 以上。先让小数据跑通,再按数据量调参,是这类项目最稳的推进方式。

4. 从特征到评分结果落库:逻辑回归评分卡与可视化看板

特征宽表生成之后,系统已经完成“大数据”的主线。接下来的模型部分不需要复杂算法,逻辑回归在这个场景里是最合适的:Spark MLlib 自带实现,训练快,权重可以直接解释,答辩时能讲清楚每个特征对风险分的影响方向。

4.1 训练一个可解释的评分卡:Spark MLlib 逻辑回归

把宽表里的数值列组装成特征向量,用标准缩放消除量纲差异,再训练逻辑回归。注意划分训练集和验证集时别用随机抽样,按时间切分更能避免“用未来数据预测过去”的穿越问题:

from pyspark.sql.functions import col from pyspark.ml.feature import VectorAssembler, StandardScaler from pyspark.ml.classification import LogisticRegression df = spark.read.parquet("/warehouse/dws/user_features") feature_cols = [ "recent_30_apply_cnt", "recent_30_amount", "overdue_cnt", "max_loan_amount", "current_loan_amount" ] df = df.withColumn( "is_overdue", (col("overdue_cnt") > 0).cast("int") ) assembler = VectorAssembler(inputCols=feature_cols, outputCol="features_vec") scaler = StandardScaler(inputCol="features_vec", outputCol="features") lr = LogisticRegression(featuresCol="features", labelCol="is_overdue") # 特征加工和标准化 df_scaled = scaler.fit(assembler.transform(df)).transform( assembler.transform(df) ) train_df = df_scaled.filter(col("reg_time") < "2020-06-01") test_df = df_scaled.filter(col("reg_time") >= "2020-06-01") model = lr.fit(train_df)

这里把is_overdue定义成“历史上有没有过逾期”,而不是“当前是否逾期”,这是评分卡里常用的二分类口径。StandardScaler 把金额类特征缩放到均值 0、方差 1,防止“借款金额”因为数值大而主导模型权重。逻辑回归训练完成后,用model.transform(test_df)就能得到每个用户的逾期概率,将概率乘以 1000 再取整,就是可视化看板上的风险分。

4.2 评分结果写 MySQL:幂等写入与分区裁剪

Spark 算完的评分结果要落到 MySQL 供后台查询,这一步看似简单,坑却不少。最常见的问题是任务失败后重跑,MySQL 里多出一批重复数据。解决办法是先按业务日期删除当天数据,再执行写入:

from pyspark.sql.functions import col, lit score_df = model.transform(test_df).select( col("user_id"), col("prediction"), lit("2024-06-01").alias("stat_date") ) # 先删后插,保证同一天跑多次不重复 spark.read.jdbc( url="jdbc:mysql://localhost:3306/risk_db", table="user_score_daily", properties={"user": "root", "password": "your_password", "driver": "com.mysql.cj.jdbc.Driver"} ).filter(col("stat_date") == "2024-06-01") \ .write.mode("overwrite") \ .jdbc(url, "user_score_daily", properties={"user": "root", "password": "your_password", "driver": "com.mysql.cj.jdbc.Driver"})

上面的写法把“查询旧数据”和“写入新数据”拼在一起,依赖 MySQL 端存在主键(user_id, stat_date)。更稳的做法是先写到临时表user_score_tmp,确认写完后用一条 SQL 把临时表数据替换进正式表,整个过程对线上服务无感。写入时记得给 spark-submit 带上 MySQL 驱动 jar,否则会报 ClassNotFound。

4.3 演示用可视化看板:最少代码出风险分分布图

风控系统通常需要一个看板展示风险分分布,毕业设计答辩时一张图比十行代码更有说服力。这里用 pyecharts 生成一个风险分区间柱状图:

from pyecharts.charts import Bar from pyecharts import options as opts df_score = spark.read.jdbc( url="jdbc:mysql://localhost:3306/risk_db", table="user_score_daily", properties={"user": "root", "password": "your_password", "driver": "com.mysql.cj.jdbc.Driver"} ) score_pd = df_score.groupBy("prediction").count().toPandas() bar = ( Bar() .add_xaxis(score_pd["prediction"].astype(str).tolist()) .add_yaxis("用户数", score_pd["count"].tolist()) .set_global_opts(title_opts=opts.TitleOpts(title="风控评分分布")) ) bar.render("score_dist.html")

toPandas()只适合结果集小的时候。看板数据量很大时,先让 Spark 聚合出每个分数区间的人数,再转 Pandas,避免把几十万行明细一次性拉进驱动节点。如果你的环境装不了 pyecharts,也可以让 Spark 把聚合结果导出成 CSV,再用 ECharts 的 HTML 模板渲染,效果一样。

5. 避坑:Hadoop 和 Spark 风控开发中的五个典型问题

跑批跑得多了,你会发现这一整套环境里翻车最多的不是算法,而是环境和数据细节。下面五条是我在这个题目上最常见的踩坑记录,每条都按现象、原因、解决三个步骤整理。

5.1 现象:HDFS Web 界面打不开,jps 里也没有 NameNode

原因分两种。一种是hdfs namenode -format只执行了一次,但格式化之后又重启了机器,NameNode 的元数据目录和 DataNode 的数据目录不一致;另一种是伪分布式下 core-site.xml 里写着hdfs://localhost:9000,但机器 hostname 配的不是 localhost,启动时绑定失败。

解决方法是检查$HADOOP_HOME/logs下的 hadoop-namenode 日志,看到Cannot connect to port 9000就去核对 fs.defaultFS。看到元数据不一致,就把/tmp/hadoop-*下的数据目录删干净,重新格式化再启动。注意格式化命令只能在第一次用,之后每次启动直接start-dfs.sh,别手滑又执行一次。

5.2 现象:Spark 作业一直卡在 ACCEPTED 状态,日志里没有任何报错

这个状态说明作业已经提交给 YARN,但迟迟没有分配到容器。原因几乎都是资源不够:executor 申请的内存或核数超过了 NodeManager 最大可用值,或者集群里其他任务占满了资源。最常见的是伪分布式机器只有 4G 内存,却配置了spark.executor.memory=4g,YARN 直接拒绝分配。

解决方法是先看 YARN Web 界面的 Cluster Metrics,确认可用内存还剩多少,然后按机器实际内存调参。单机伪分布式我用--num-executors 1 --executor-memory 1g --executor-cores 1,几乎不会卡资源。如果任务必须跑大内存,就调大yarn.scheduler.maximum-allocation-mb,这个值默认只有 1G,经常是罪魁祸首。

5.3 现象:特征宽表 join 完之后行数比用户数还多

宽表的目标是每个用户一行,但 left join 之后出现了大量重复用户。原因是在聚合之前做了多表 join:用户表先 join 借款明细,又 join 还款明细,两个一对多关系叠在一起,行数变成了笛卡尔积。这不是 Spark 的 bug,是关联逻辑错了。

解决方法是所有明细表先自己完成 groupBy 聚合,把每个用户压成一行,再和用户表 left join。我在 3.3 节里就是这么写的。判重的方法也简单:join 后执行df_wide.groupBy("user_id").count().filter("count > 1"),如果结果为空,宽表才合格。

5.4 现象:MySQL 里出现重复评分数据,同一用户同一天有两条记录

这个问题的根源通常不在代码,而在调度:任务失败后重跑,写入时没有清旧数据。第一次写进 10 万条,失败重跑又加了 10 万条,或者写入时用了append模式。解决方法是给目标表加联合主键(user_id, stat_date),每次写入前先按 stat_date 删除历史分区数据,或者走临时表替换。你可以跟答辩老师直接说:这里做过幂等设计,保证批任务任意重跑结果一致。

5.5 现象:跑批时间越来越长,从 10 分钟慢慢变成 30 分钟

这是典型的“小文件问题”。每次 Spark 作业写 Parquet 默认会按 shuffle 分区生成文件,几十个分区就是几十个小文件,长期积累后 HDFS 上挤满了 KB 级文件,NameNode 压力大,读取时任务数暴增。解决思路是给写出的 DataFrame 做一次repartition()控制文件数量,或者在写 Hive 表时按日期分区,每个分区落地成一个较大的文件。

df_wide.coalesce(1) \ .write.mode("overwrite") \ .parquet("/warehouse/dws/user_features")

coalesce(1)会把结果合并成一个大文件,但数据量大时也降低了后续读取的并行度,所以更通用的做法是repartition(50),让文件数量和集群并发度匹配。日志里如果看到大量 job 都在处理几百 KB 的输入,基本就是小文件在作祟。

6. 让毕设从“能跑”升级到“能答辩”:性能验证与演示细节

系统能跑只是及格线,答辩时需要用数据证明“Hadoop 和 Spark 的选型是有必要性的”。我建议你准备一个数据量对比实验:同一份特征加工代码,分别在 5 万行和 200 万行数据上各跑一次,记录耗时和资源占用。结果通常是小数据量下 Spark 因为有任务调度开销反而比 Pandas 慢,但数据量越大,Spark 和单机框架的差距越明显。这个实验不用搞成严格的性能报告,只要在论文里放一张耗时对比表,就能支撑“为什么需要 Spark”这个必考题。

演示顺序也有讲究,我习惯按三条线走:先讲数据流向,从 HDFS 原始文件到 Parquet 宽表,展示 Spark SQL 的执行计划;再讲任务流向,展示 YARN 上 Spark Application 的提交过程、executor 的分配情况;最后讲异常处理,直接现场看一遍日志里某个失败任务的报错和重跑结果。前两条证明你理解了系统,第三条证明你有排查能力,这比背代码更让老师信服。

最后提醒一个我在答辩时被当场追问过的问题:为什么用的是批处理而不是 Spark Streaming。这个问题没有标准答案,但你必须能说清楚边界——信贷风控的评分多数场景是 T+1 离线批处理,实时计算通常只用于反欺诈拦截,那是另一条技术栈。把离线链路讲透,比硬塞一个实时模块更稳妥。这个方向做完,你手里的不仅是一份能跑的代码,还是一个能讲清楚选型理由、踩过真实坑位的完整项目。希望帮到你。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询