简介:这份资源是一套基于Hadoop与Spark的金融信贷风险大数据分析系统源码,面向计算机相关专业的毕业设计学生及需要项目实训的学习者,可作为课程设计或期末综合作业的实践材料。系统以Hadoop负责海量数据存储与批处理、Spark提供内存计算能力,实现从数据采集、清洗到特征提取、模型训练的全流程风控分析,涵盖用户管理、数据接入、风险指标计算与可视化展示等核心模块。压缩包共83个文件,约64KB,以36个Java源码、8个Scala程序、12个XML配置、5个properties配置及SQL脚本、说明文档等为主,另含备份文件与前端页面资源,结构完整、层次清晰。已有104人学习关注。项目经导师审核并获优秀评级,代码经过系统化调试,各组件接口衔接完整,部署说明清晰,读者可据此理解分布式计算在金融风控领域的落地方法,掌握风险控制系统的构建原理,并在此基础上进行扩展开发。
1. 从一份信贷风控需求说起:Hadoop 与 Spark 到底在这套系统里干什么
去年帮一个做消费金融的朋友看他们新上的风控报表,业务方要的是「每个申请人在放款前,能看到他近 6 个月在多头平台的申贷次数、当前负债率、以及同设备关联的逾期率」。数据量不算夸张,日均申请件 40 万左右,但麻烦在于数据散在 MySQL 业务库、埋点日志、第三方征信回执文件里,用单机 Python 跑一次全量要 6 个多小时,跑完当天的审批高峰早过了。后来把离线特征计算挪到 Hadoop + Spark 上,同样的逻辑压到 20 分钟以内,才有了「T+1 早上出特征、白天审批直接查」的节奏。
这套「基于 Hadoop 与 Spark 的金融信贷风险大数据分析系统」,本质就是解决上面这类问题:用 HDFS 把多源信贷数据沉下来,用 Spark 做清洗、特征加工和风险指标计算,最后把结果落到能对接审批流的存储里。它适合两类人——一类是正在做大数据方向毕业设计、需要一套能跑通、能讲清链路的完整方案;另一类是刚接手信贷风控数据、想用离线批处理替代单机脚本的工程师。下面按「数据怎么进、特征怎么算、坑在哪」的顺序拆开讲,代码和参数都能直接抄。
2. 数据链路与存储选型:HDFS 分层、Hive 建表与信贷字段设计
2.1 为什么信贷风控场景优先选 HDFS + Hive 而不是直接上数据库
信贷数据的第一个特点是「写一次、读多次、很少改」。申请件一旦落库,除了少量状态回写,绝大多数场景是拿历史数据反复算特征。这种读多写少、批量扫描的模式,正好是 HDFS 的强项。第二个特点是字段会膨胀——今天算 30 个特征,下个月业务要加「近 3 个月夜间申请占比」,如果一开始用关系库,加字段、改表结构、补历史数据都很别扭;用 Hive 外部表 + Parquet,加列基本不影响存量数据。
常见做法是把数据按「原始层 / 清洗层 / 特征层」三层放。原始层保留第三方回执的原始 JSON 和业务库的 binlog 落地文件,清洗层做去重、脱敏、类型统一,特征层才是 Spark 计算后给审批用的宽表。这样做的直接好处是:特征算错了,能回溯到清洗层甚至原始层重跑,不用求业务库重新导数据。
提示:毕业设计里如果只搭了伪分布式,三层目录照样要建,评审看的是分层思路,不是集群规模。
2.2 HDFS 目录规划与 Hive 建表语句
先建目录。下面这段在 HDFS 上按日期分区建三层路径,${bizdate}用调度传进来的业务日期替换:
# 原始层:第三方回执、埋点日志原样落地 hdfs dfs -mkdir -p /credit/ods/receipt/dt=${bizdate} hdfs dfs -mkdir -p /credit/ods/apply_log/dt=${bizdate} # 清洗层:去重脱敏后的明细 hdfs dfs -mkdir -p /credit/dwd/apply_detail/dt=${bizdate} # 特征层:给审批用的宽表 hdfs dfs -mkdir -p /credit/dws/risk_feature/dt=${bizdate}逻辑说明:按dt=分区是最省事的做法,Spark 写的时候用partitionBy("dt")自动对齐,后续查某一天只扫一个目录。参数上,bizdate建议统一用yyyyMMdd,别混用yyyy-MM-dd,否则 Hive 分区字段类型和路径对不上,查出来是空。
清洗层建表,字段围绕信贷风控最常用的几类:申请人标识、申请时间、渠道、设备、额度、以及第三方返回的多头指标。
CREATE EXTERNAL TABLE IF NOT EXISTS credit_dwd.apply_detail ( apply_id STRING COMMENT '申请单号', cust_id STRING COMMENT '客户号(脱敏)', id_md5 STRING COMMENT '身份证MD5', phone_md5 STRING COMMENT '手机号MD5', device_id STRING COMMENT '设备指纹', channel STRING COMMENT '进件渠道', apply_time STRING COMMENT '申请时间', apply_amount DECIMAL(12,2) COMMENT '申请金额', multi_loan_cnt INT COMMENT '近6月多头申贷次数', overdue_flag INT COMMENT '当前是否逾期 0/1' ) COMMENT '信贷申请清洗明细' PARTITIONED BY (dt STRING) STORED AS PARQUET LOCATION '/credit/dwd/apply_detail';逻辑说明:用EXTERNAL是为了删表不删数据,重跑时只覆盖分区。id_md5、phone_md5存的是脱敏后的值,原始证件号不进数仓,这是信贷场景的硬要求。multi_loan_cnt这类第三方指标直接平铺成列,比塞进 JSON 字符串再解析要快得多,代价是加指标要改表——但信贷指标相对稳定,这个取舍划算。
参数上,DECIMAL(12,2)对应金额,别用DOUBLE,否则对账时会出现0.30000000000000004这种尾差,风控对账最忌讳这个。分区字段dt放最后,Hive 会自动识别。
2.3 数据接入的两种常见方式
业务库数据一般走 Sqoop 或 DataX 抽到 HDFS,第三方回执文件走hdfs dfs -put或 Flume。毕业设计环境里,我一般建议直接用 Spark 读本地 CSV/JSON 再写 HDFS,省掉 Sqoop 配置的麻烦:
# 把本地第三方回执 JSON 读进来,落到清洗层 df = spark.read.json("file:///data/receipt/20240501/*.json") df.write.mode("overwrite").parquet("/credit/ods/receipt/dt=20240501")逻辑说明:file:///前缀是本地路径,不加会被当成 HDFS 路径。mode("overwrite")保证重跑幂等,同一天跑两次不会翻倍。这一步只做落地,不做清洗,清洗放到下一章,职责分开,出问题好定位。
3. Spark 特征计算实战:多头指标、负债率与逾期关联的算子写法
3.1 特征计算为什么用 Spark SQL 而不是 RDD
信贷特征绝大多数是「分组聚合 + 窗口」的形态:按客户分组算近 6 个月申贷次数、按设备分组算关联逾期率、按时间窗口算夜间申请占比。这些用 Spark SQL 的窗口函数写,比 RDD 的groupByKey清晰得多,而且 Catalyst 优化器会自动做谓词下推和列裁剪,扫 Parquet 时只读用到的列。血泪经验是:早期用 RDD 手写聚合,同样的逻辑代码量翻倍,还容易在 shuffle 阶段把内存打爆。
常见做法是把特征拆成几个独立的 SQL 或 DataFrame 任务,各自算完再 join 成宽表。别写一个巨型 SQL 从头算到尾,中间结果不落盘,一旦某步出错整条链重跑,调试成本极高。
3.2 多头申贷次数与负债率计算
下面这段算两个核心特征:近 6 个月多头申贷次数、当前负债率。负债率用「已用额度 / 授信总额」近似,实际项目里授信总额来自征信回执。
from pyspark.sql import SparkSession, functions as F from pyspark.sql.window import Window spark = SparkSession.builder \ .appName("credit_risk_feature") \ .config("spark.sql.shuffle.partitions", "200") \ .enableHiveSupport() \ .getOrCreate() detail = spark.table("credit_dwd.apply_detail").filter("dt = '20240501'") # 近6月多头申贷次数:按客户+身份证聚合,统计不同机构数 multi = detail.groupBy("cust_id", "id_md5") \ .agg(F.countDistinct("channel").alias("multi_loan_cnt"), F.sum("apply_amount").alias("total_apply_amt")) # 负债率:已用额度/授信总额,授信来自征信字段 debt = detail.groupBy("cust_id") \ .agg((F.sum("used_amount") / F.sum("credit_limit")).alias("debt_ratio")) # 按申请时间做窗口,算近30天申请频次 w = Window.partitionBy("cust_id").orderBy(F.col("apply_time").cast("long")) \ .rangeBetween(-30 * 86400, 0) freq = detail.withColumn("apply_cnt_30d", F.count("apply_id").over(w)) \ .select("cust_id", "apply_id", "apply_cnt_30d") feature = multi.join(debt, "cust_id", "left") \ .join(freq, "cust_id", "left") \ .fillna({"debt_ratio": 0.0, "apply_cnt_30d": 0}) feature.write.mode("overwrite").partitionBy("dt") \ .parquet("/credit/dws/risk_feature")逻辑说明:countDistinct("channel")统计的是不同进件渠道数,用来近似多头机构数——真实项目里应该用机构编码字段,这里用渠道代替是为了演示。rangeBetween(-30*86400, 0)是按秒做时间窗口,apply_time必须先转成时间戳,否则窗口按字符串排序会乱。fillna把没匹配上的负债率补 0,避免下游审批因为 null 直接拒件。
参数上,spark.sql.shuffle.partitions默认 200,小数据量下反而慢,毕业设计环境可以调到 20~50;生产环境按数据量调,一般每个分区 128MB 左右。enableHiveSupport()是为了能直接spark.table读 Hive 表,不加就得写全路径。
3.3 逾期关联与设备维度特征
设备维度是信贷反欺诈的重点:同一设备关联多个申请人、且这些申请人有逾期,就是高风险信号。
# 设备关联逾期率 device_risk = detail.groupBy("device_id") \ .agg(F.countDistinct("cust_id").alias("device_cust_cnt"), F.avg("overdue_flag").alias("device_overdue_rate")) \ .filter("device_cust_cnt >= 3") # 至少关联3个客户才算团伙 feature2 = feature.join(device_risk, "device_id", "left") \ .fillna({"device_overdue_rate": 0.0})逻辑说明:filter("device_cust_cnt >= 3")是业务阈值,关联客户太少的设备统计意义不大,反而引入噪声。device_overdue_rate是设备下所有客户逾期标记的均值,直接当风险分用。这一步 join 后要检查数据量有没有膨胀——如果device_risk里一个设备对应多行,join 会炸,所以前面必须groupBy("device_id")保证唯一。
注意:join 前一定先
count一下两边行数,信贷数据里设备 ID 为空或重复的情况很常见,join 键不唯一是这类任务最常见的翻车点。
4. 避坑与排查:信贷大数据任务里最容易翻车的五件事
4.1 数据倾斜导致个别 task 卡死
现象:Spark 任务 199 个 task 秒完,剩 1 个跑 40 分钟,日志里某个 key 的记录数是其他 key 的几百倍。原因:信贷数据里「默认设备」「空手机号」这类脏值会聚成超大 key,groupBy时全压到一个分区。解决:先过滤脏值,或对热点 key 加随机前缀打散再聚合:
# 对空设备ID加随机后缀打散 detail = detail.withColumn("device_id", F.when(F.col("device_id").isNull() | (F.col("device_id") == ""), F.concat(F.lit("unknown_"), (F.rand() * 10).cast("int"))) .otherwise(F.col("device_id")))4.2 时间字段类型不统一导致窗口算错
现象:apply_cnt_30d算出来全是 1,或者窗口范围明显不对。原因:apply_time在原始数据里是字符串2024-05-01 10:00:00,直接cast("long")得到的是 null 或错误值。解决:统一用to_timestamp转换,再转时间戳:
detail = detail.withColumn("apply_ts", F.unix_timestamp(F.to_timestamp("apply_time", "yyyy-MM-dd HH:mm:ss")))4.3 分区字段写错导致查不到数据
现象:任务跑成功,但select * from ... where dt='20240501'返回空。原因:写的时候partitionBy("dt")但 DataFrame 里根本没有dt列,Spark 会报错或写出无分区目录。解决:写之前显式加分区列:
feature = feature.withColumn("dt", F.lit("20240501")) feature.write.partitionBy("dt").parquet("/credit/dws/risk_feature")4.4 小文件过多拖垮 NameNode
现象:特征层目录下几万个几十 KB 的小文件,下次读的时候启动就卡。原因:分区太细 + 并行度太高,每个 task 写一个文件。解决:写完做一次合并,或写之前repartition:
feature.repartition(10).write.mode("overwrite") \ .partitionBy("dt").parquet("/credit/dws/risk_feature")4.5 内存不足报 OOM 却以为是数据量问题
现象:java.lang.OutOfMemoryError,但数据量并不大。原因:多半是collect()或toPandas()把大结果拉回 Driver,或者 join 时广播了不该广播的大表。解决:去掉 Driver 端收集,检查 broadcast join 阈值:
# 关掉自动广播,避免大表被广播 spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")提示:信贷特征任务里,凡是
collect、toPandas、show大结果集的操作,上线前都要删掉,这是最常见的 Driver OOM 来源。
5. 从跑通到可信:特征校验、调度与一个我每次都走的验证习惯
特征算出来不等于能用。信贷场景里,一个算错的多头指标可能直接导致误拒,所以上线前必须做校验。我一般固定走三步:分布校验、抽样对账、边界值检查。
分布校验是看特征值的分布有没有突变。比如multi_loan_cnt昨天均值 2.3,今天突然 8.7,大概率是上游数据源换了或去重逻辑坏了。用 Spark 一行就能出:
feature.select( F.avg("multi_loan_cnt").alias("avg_multi"), F.expr("percentile_approx(multi_loan_cnt, 0.5)").alias("p50_multi"), F.expr("percentile_approx(multi_loan_cnt, 0.99)").alias("p99_multi"), F.sum(F.when(F.col("multi_loan_cnt").isNull(), 1).otherwise(0)).alias("null_cnt") ).show()逻辑说明:percentile_approx比精确分位数快得多,风控看分布够用。null_cnt单独统计,因为 null 在审批里可能被当成 0,掩盖了数据缺失。参数上,0.99分位用来抓极端值,如果 p99 突然从 5 跳到 50,基本可以断定有脏数据混进来了。
抽样对账是拿 100 个客户,用单机 Python 从原始数据重算一遍,和 Spark 结果逐条比。这一步最笨但最有效,我见过太多「逻辑看着对、结果差一点」的情况,都是 join 键重复或时间窗口边界没对齐导致的。对账脚本不用写得多优雅,能跑出差异清单就行。
边界值检查针对的是空值、负数、超大值。信贷金额出现负数、负债率大于 1、申请次数为负,这些都要在写特征层之前拦掉:
feature = feature.filter( (F.col("debt_ratio").between(0, 1)) & (F.col("multi_loan_cnt") >= 0) & (F.col("apply_amount") > 0) )调度上,毕业设计环境用 crontab 或 Airflow 单机版都行,关键是任务要有依赖:清洗层跑完才能跑特征层,特征层跑完才能跑校验。别把三个任务写成一个脚本从头跑到尾,中间任何一步失败都得全量重来。生产环境常见做法是 Airflow 的ExternalTaskSensor或 DolphinScheduler 的依赖节点,原理一样。
最后说个我自己的习惯:每次改完特征逻辑,不管改动多小,都强制走一遍「小样本跑通 → 分布对比 → 抽样对账」这三步,哪怕只是改了个fillna的值。有一次就是把fillna(0)改成fillna(0.0),看着没区别,结果下游按字符串解析时把0.0当成了非法值,整批特征被拒。从那以后我每次上线前都强制走一遍这三步,再急也不跳。希望这套链路和踩坑记录能帮到你,把这份方案真正跑起来、用起来。
本文还有配套的精品资源,点击获取