简介:这份资源是面向计算机相关专业学生与开发者的毕业设计项目源码,主题为基于Hadoop与Spark的大数据金融信贷风险控制系统,适合用作毕设、课程设计、作业或项目初期立项演示,也便于基础较好的学习者在此基础上二次修改扩展功能。压缩包共23个文件,约53KB,以Scala核心代码与XML配置为主,辅以properties参数文件、iml工程标识、js前端脚本及json数据文件,整体结构紧凑,涵盖数据源接入、Spark流式处理与信贷风控业务逻辑等模块,并附有README说明文档,便于快速理解项目组织方式。目前已有2379人浏览学习,代码经过运行测试,功能可正常使用,下载后可直接参考其目录划分与实现思路,用于梳理大数据风控项目的技术栈与开发流程。
1. 从一份毕业设计源码说起:Hadoop+Spark 信贷风控系统到底在做什么
信贷风控这个场景,很多人第一反应是“评分卡”和“规则引擎”,但真正让一套系统能扛住日均几十万笔申请、几千万条行为埋点的,是底层那套分布式计算链路。这份毕业设计标题里的 Hadoop+Spark 组合,本质上解决的是同一件事:把散落在业务库、日志文件、第三方征信接口里的异构数据,用 HDFS 统一存下来,再用 Spark 做特征加工、规则匹配和模型打分,最后把结果写回业务库供审批调用。它适合两类人:一类是正在做大数据方向毕业设计、需要一套能跑通全流程参考实现的学生;另一类是刚转行到金融科技公司、想搞清楚离线风控批处理链路怎么搭的初中级工程师。源码本身不是重点,重点是它背后那条“数据接入→清洗→特征→规则→打分→回写”的链路,你能不能在自己的机器上复现出来,以及复现过程中那些没人告诉你的参数和坑。
2. 环境先立住:Hadoop 伪分布式与 Spark 本地模式的最小可用组合
2.1 为什么毕业设计场景优先选伪分布式而不是全分布式
全分布式集群听起来更“大数据”,但对毕业设计来说,三台虚拟机带来的网络配置、时间同步、SSH 免密、防火墙策略问题,会吃掉你至少两天时间,而这两天你本可以用来调通风控规则。伪分布式在一台机器上模拟 NameNode、DataNode、ResourceManager、NodeManager 全部角色,HDFS 和 YARN 的 API 行为与全分布式一致,Spark 提交任务时看到的资源调度逻辑也相同。唯一区别是数据块副本数只能设为 1,吞吐量上不去,但风控系统的离线批处理对吞吐不敏感,对逻辑正确性敏感。常见做法是:Hadoop 用 3.3.x 稳定版,Spark 用 3.4.x 或 3.5.x,JDK 锁死 1.8,Scala 用 2.12。版本不要追新,Spark 3.5 对 JDK 17 的支持在伪分布式下仍有少量序列化问题,毕业设计求稳不求新。
2.2 从零开始安装 Hadoop 的关键命令与配置项
先确认机器有至少 8GB 内存、50GB 空闲磁盘,然后按下面步骤走。不要跳过任何一步,尤其是环境变量和配置文件里的路径。
# 创建专用用户,避免权限混乱 sudo useradd -m hadoop sudo passwd hadoop su - hadoop # 下载并解压 Hadoop 3.3.6 wget https://archive.apache.org/dist/hadoop/common/hadoop-3.3.6/hadoop-3.3.6.tar.gz tar -zxvf hadoop-3.3.6.tar.gz -C /home/hadoop/ mv /home/hadoop/hadoop-3.3.6 /home/hadoop/hadoop # 配置环境变量 echo 'export HADOOP_HOME=/home/hadoop/hadoop' >> ~/.bashrc echo 'export PATH=$PATH:$HADOOP_HOME/bin:$HADOOP_HOME/sbin' >> ~/.bashrc echo 'export JAVA_HOME=/usr/lib/jvm/java-8-openjdk-amd64' >> ~/.bashrc source ~/.bashrc逻辑说明:Hadoop 的启动脚本依赖JAVA_HOME和HADOOP_HOME两个变量,缺一个就会报Error: JAVA_HOME is not set。参数上,-C指定解压目录,mv重命名是为了后续配置文件路径统一。接下来改四个核心配置文件,路径在$HADOOP_HOME/etc/hadoop/下。
core-site.xml里设fs.defaultFS为hdfs://localhost:9000,这是 HDFS 的入口地址,Spark 读写 HDFS 时靠它定位。hdfs-site.xml里设dfs.replication为 1,伪分布式下副本数超过 1 会一直报副本不足。mapred-site.xml设mapreduce.framework.name为yarn,让 MapReduce 作业跑在 YARN 上。yarn-site.xml设yarn.nodemanager.aux-services为mapreduce_shuffle,这是 NodeManager 启动 shuffle 服务的必要配置。改完后执行hdfs namenode -format格式化,再start-dfs.sh和start-yarn.sh。用jps看到 NameNode、DataNode、ResourceManager、NodeManager、SecondaryNameNode 五个进程才算成功。
2.3 Spark 本地模式跑通风控特征加工的最小验证
Spark 装好后不需要改太多配置,重点是让spark-shell能读到 HDFS 上的文件。把spark-env.sh里的JAVA_HOME和HADOOP_CONF_DIR指向正确路径,然后启动。
tar -zxvf spark-3.5.0-bin-hadoop3.tgz -C /home/hadoop/ mv /home/hadoop/spark-3.5.0-bin-hadoop3 /home/hadoop/spark echo 'export SPARK_HOME=/home/hadoop/spark' >> ~/.bashrc echo 'export PATH=$PATH:$SPARK_HOME/bin' >> ~/.bashrc source ~/.bashrc # 启动 spark-shell,指定 master 为 local[2] spark-shell --master local[2]进入 shell 后跑一段最小验证,从 HDFS 读一个 CSV 并统计行数:
// 假设 HDFS 上已有 /user/hadoop/credit_apply.csv val df = spark.read.option("header", "true").csv("hdfs://localhost:9000/user/hadoop/credit_apply.csv") df.count() df.printSchema()逻辑说明:local[2]表示用两个线程模拟集群,适合验证逻辑。option("header", "true")让第一行作为列名,不设的话所有列名会变成_c0、_c1。printSchema用来确认字段类型,风控数据里金额字段如果被推断成 string,后续做聚合会报类型错误,需要显式cast("double")。这一步跑通,说明 Hadoop 和 Spark 的联通没问题,可以开始灌真实数据了。
3. 信贷风控核心链路:从原始申请数据到风险评分的 Spark 实现
3.1 数据接入层:把 CSV、JSON、MySQL 三种来源统一成 DataFrame
真实风控系统的数据来源至少三类:业务库导出的申请信息 CSV、第三方征信返回的 JSON、内部行为埋点的 MySQL 表。Spark 的 DataFrame API 能把它们统一成同一套抽象,后续特征加工不用关心底层格式。下面是一个典型的接入代码。
from pyspark.sql import SparkSession from pyspark.sql.functions import col, from_json from pyspark.sql.types import StructType, StructField, StringType, DoubleType spark = SparkSession.builder \ .appName("CreditRiskETL") \ .master("local[2]") \ .config("spark.sql.shuffle.partitions", "4") \ .getOrCreate() # 1. 读 CSV 申请数据 apply_df = spark.read.option("header", "true").option("inferSchema", "true") \ .csv("hdfs://localhost:9000/user/hadoop/apply.csv") # 2. 读 JSON 征信数据,显式定义 schema 避免推断错误 credit_schema = StructType([ StructField("apply_id", StringType(), True), StructField("credit_score", DoubleType(), True), StructField("overdue_count", DoubleType(), True) ]) credit_df = spark.read.schema(credit_schema).json("hdfs://localhost:9000/user/hadoop/credit.json") # 3. 读 MySQL 行为数据 behavior_df = spark.read.format("jdbc") \ .option("url", "jdbc:mysql://localhost:3306/risk") \ .option("dbtable", "user_behavior") \ .option("user", "root") \ .option("password", "123456") \ .load() # 4. 按 apply_id 关联 joined_df = apply_df.join(credit_df, on="apply_id", how="left") \ .join(behavior_df, on="apply_id", how="left") joined_df.cache()逻辑说明:spark.sql.shuffle.partitions默认是 200,本地模式下调成 4 能减少小文件问题。JSON 读取必须显式给 schema,因为征信接口返回的字段类型不稳定,inferSchema可能把overdue_count推断成 string。cache()是因为后续要多次用到这个关联结果,不缓存的话每次 action 都会重新读三个源。参数上,JDBC 的dbtable可以写表名也可以写子查询,写子查询时要用括号包起来并加别名。
3.2 特征加工:逾期率、收入负债比、多头申请次数的计算逻辑
风控特征分三类:稳定性特征、负债类特征、行为类特征。下面这段代码计算三个最核心的指标。
from pyspark.sql.functions import when, count, sum, avg, datediff, current_date # 逾期率:历史逾期次数 / 总还款期数 feature_df = joined_df.withColumn( "overdue_rate", when(col("total_periods") > 0, col("overdue_count") / col("total_periods")).otherwise(0.0) ) # 收入负债比:月还款额 / 月收入 feature_df = feature_df.withColumn( "dti", when(col("monthly_income") > 0, col("monthly_payment") / col("monthly_income")).otherwise(999.0) ) # 近 30 天多头申请次数 recent_apply = behavior_df.filter( datediff(current_date(), col("apply_date")) <= 30 ).groupBy("user_id").agg(count("apply_id").alias("apply_count_30d")) feature_df = feature_df.join(recent_apply, on="user_id", how="left") \ .fillna({"apply_count_30d": 0})逻辑说明:when(...).otherwise(...)是 Spark 里的条件分支,比 UDF 快很多。dti在收入为 0 时给 999 是一个业务约定,表示“无法计算,风险极高”,不要给 0,否则会被模型当成低风险。datediff返回天数差,current_date()是当前日期,这个过滤条件在批处理里每天跑一次,等价于滚动窗口。fillna把没有行为的用户补 0,不补的话后续聚合会出 null,导致评分卡分箱报错。
3.3 规则引擎与评分卡:用 Spark SQL 实现可配置的风险决策
规则引擎的核心是把业务规则从代码里抽出来,变成一张配置表,Spark 读表后动态生成 SQL 条件。下面是一个简化实现。
# 规则配置表结构:rule_id, rule_name, condition_sql, score, decision rules = [ ("R001", "高负债比", "dti > 0.7", -50, "REJECT"), ("R002", "近期多头", "apply_count_30d > 5", -30, "REVIEW"), ("R003", "历史逾期", "overdue_rate > 0.1", -40, "REJECT"), ("R004", "征信分低", "credit_score < 600", -20, "REVIEW") ] # 动态拼接评分表达式 score_expr = "0" for rule_id, name, cond, score, decision in rules: score_expr = f"CASE WHEN {cond} THEN {score_expr} + ({score}) ELSE {score_expr} END" result_df = feature_df.selectExpr("*", f"{score_expr} as risk_score") # 根据总分和单条规则决定最终决策 result_df = result_df.withColumn( "final_decision", when(col("risk_score") <= -80, "REJECT") .when(col("risk_score") <= -40, "REVIEW") .otherwise("PASS") ) result_df.select("apply_id", "risk_score", "final_decision") \ .write.mode("overwrite").parquet("hdfs://localhost:9000/user/hadoop/risk_result")逻辑说明:CASE WHEN嵌套是 Spark SQL 里做多条件累加的标准写法,比写多个withColumn再相加更清晰。规则表在真实系统里存在 MySQL 或 ZooKeeper 里,这里为了演示直接写在代码里。write.mode("overwrite")每天覆盖当天分区,生产环境一般按日期分区写partitionBy("dt")。参数上,risk_score的阈值 -80 和 -40 是业务调出来的,不是拍脑袋定的,需要根据通过率反推。
4. 避坑与排查:信贷风控 Spark 作业最常见的五类翻车
4.1 现象:任务卡在 99% 不动,日志显示 shuffle fetch failed
原因:伪分布式下yarn.nodemanager.resource.memory-mb默认值太小,Spark executor 申请不到足够内存,shuffle 阶段数据拉取超时。解决:在yarn-site.xml里把该值调到 4096 以上,同时在 Spark 提交时加--executor-memory 2g --driver-memory 2g。如果还不行,检查/etc/hosts里 localhost 是否解析到 127.0.0.1,有些系统会解析到 IPv6 地址导致连接失败。
4.2 现象:DataFrame 关联后行数暴涨,从 10 万变成 50 万
原因:关联键有重复值,比如apply_id在征信表里一个申请对应多条征信查询记录。解决:关联前先对征信表做groupBy("apply_id").agg(max("credit_score"))去重,或者改用窗口函数取最新一条。血泪经验是:任何 join 之前先用df.groupBy(key).count().filter("count > 1").show()检查键唯一性,这个习惯能省掉大量排查时间。
4.3 现象:写入 HDFS 时报FileAlreadyExistsException
原因:Spark 的write默认模式是errorifexists,目标路径已存在就报错。解决:改成mode("overwrite")或mode("append")。但注意,overwrite会删掉整个目录再重建,如果目录下有其他分区数据会一起丢。生产环境更安全的做法是按日期分区写,每次只覆盖当天分区。
4.4 现象:中文乱码,姓名和地址字段全是问号
原因:CSV 文件编码是 GBK,Spark 默认按 UTF-8 读。解决:读取时加.option("encoding", "GBK"),或者提前用iconv -f GBK -t UTF-8转码。注意,encoding选项在 Spark 3.x 的 CSV reader 里支持,JSON reader 不支持,JSON 必须提前转码。
4.5 现象:评分结果每天波动很大,同一用户昨天 PASS 今天 REJECT
原因:current_date()在批处理里取的是服务器时间,如果任务跨天跑,前后两批数据的时间基准不一致。解决:把跑批日期作为参数传进来,用lit(run_date)替代current_date(),保证同一次跑批内所有时间计算基于同一个日期。这个坑在毕业设计答辩时经常被老师问到,提前改掉能加分。
5. 进阶技巧:用 Spark 内存监测和分区调优把跑批时间压下来
5.1 用 Spark UI 定位内存瓶颈
Spark 提交任务后,默认在 4040 端口开 Web UI。重点看两个地方:Stage 页面的 Shuffle Read/Write 大小,以及 Executor 页面的 Storage Memory 使用率。如果某个 Stage 的 Shuffle Write 超过 1GB,说明分区数太少,需要调大spark.sql.shuffle.partitions。如果 Storage Memory 长期在 90% 以上,说明cache()的数据太多,要改成persist(StorageLevel.MEMORY_AND_DISK)。我一般会在代码里加一行spark.sparkContext.setLogLevel("WARN"),减少日志干扰,然后盯着 UI 看哪个 Stage 最慢。
5.2 分区调优:从 200 到 4 再到动态调整
本地模式下spark.sql.shuffle.partitions设 4 就够,但数据量上到百万级后,4 个分区每个要处理 25 万条,容易 OOM。一个实用的估算公式是:分区数 = 数据量(MB) / 128。比如 500MB 的 shuffle 数据,设 4 个分区,每个分区 125MB,接近 128MB 的 HDFS 块大小,比较合理。如果数据量波动大,可以用 Adaptive Query Execution(AQE),在 Spark 3.x 里设spark.sql.adaptive.enabled=true,让 Spark 在运行时自动合并小分区。
5.3 一个具体的验证方法:对比调优前后的跑批耗时
不要凭感觉说“调优了”,要拿数据说话。下面这段脚本记录跑批时间并写入日志。
import time start = time.time() result_df.write.mode("overwrite").parquet("hdfs://localhost:9000/user/hadoop/risk_result") end = time.time() with open("/home/hadoop/spark_job_time.log", "a") as f: f.write(f"run_date={run_date}, elapsed={end-start:.2f}s, partitions={spark.conf.get('spark.sql.shuffle.partitions')}\n")逻辑说明:time.time()返回秒级浮点数,elapsed是端到端耗时,包含 Spark 的懒执行触发时间。spark.conf.get读取当前分区配置,方便对比不同参数下的耗时。我一般会跑三组:分区数 4、8、16,各跑三次取平均,然后选耗时最低且没有 OOM 的那组。这个习惯是从一次线上事故里学来的——当时凭经验设了 200,结果本地模式启动 200 个任务,光调度开销就花了 40 秒。希望帮到你。
本文还有配套的精品资源,点击获取