
简介面向大数据金融信贷风控领域学习者和毕业设计开发者的完整项目源码包基于Hadoop与Spark技术栈实现信贷风险控制系统覆盖数据接入、流式处理、风控逻辑及可视化等环节适合课程设计、毕设或项目初期演示。压缩包内共69个文件以36个Java源码、8个Scala源码为主配合12个XML配置、5个properties属性文件、SQL脚本及H5前端JS等整体约58KB工程结构清晰便于按功能模块查阅。已有325人学习下载代码均测试运行成功答辩平均分达96分并提供远程教学支持。参考README与数据库脚本即可快速搭建环境重点理解Spark Streaming实时数据接入、MyBatis映射配置以及风控流程设计等实践要点下载后可直接导入IDE运行调试便于二次开发。1. 基于Hadoop、Spark的信贷风控系统在解决什么先看清它要扛下的数据量做信贷业务的同学都熟悉这个画面一天涌入几十万笔申请每笔申请背后是用户基本信息、多头借贷记录、设备指纹、APP行为轨迹、第三方黑名单光原始字段就上百个。传统MySQL单库撑死几千QPS十几张表做一次join直接超时更别提把一年甚至三年的历史数据拉出来做特征回溯。基于Hadoop、Spark的大数据金融信贷风控系统核心就是把数据存储和特征计算搬到分布式框架上HDFS承接海量明细数据Spark跑离线特征工程、批量评分和模型训练再配合规则引擎处理实时审批决策。它解决的是三个具体问题——数据存不下、特征算不动、风险评分出不来。适合正在做信贷技术选型的工程师也适合高校大数据方向拿一个完整项目做课程设计的同学。你先搞清楚这套系统的数据链路和关键取舍再去看源码和文档效率会高很多。2. 架构设计与技术选型为什么是Hadoop加Spark而不是其他组合2.1 先做减法哪些方案被排除理由是什么常见的第一反应是用MySQL加定时任务搞定一切。但在信贷风控场景里原始数据里有用户授权读取的运营商通话详单、电商消费记录、社保公积金流水单用户每月产生的明细记录就有几百条百万用户就是几亿行。MySQL单表过亿之后即使加了索引复杂的聚合查询也要几十秒甚至分钟级根本撑不住风控调参时反反复复的特征回溯。纯Flink流处理方案也是一个选择但风控场景里真正耗费算力的不是实时计算本身而是全量数据的批量特征计算、模型训练和离线回测。Flink擅长的是秒级窗口内的实时计算把这部分任务硬塞给Flink集群成本会翻很多倍而且Flink对状态后端的管理也比Spark的批处理模型复杂。Hadoop加Spark是更成熟的配合方式HDFS做分布式存储底座Spark跑内存计算MapReduce留作极重的历史回溯兜底任务。这套组合在金融行业被验证了超过十年招聘市场上也确实大量岗位要求这两项技术。2.2 总体架构五层数据流和一个核心整套系统的数据流可以拆成五个层级采集层Flume采集服务器日志Kafka承接业务系统实时推送的申请事件两份数据最终都落到HDFS。存储层HDFS存放原始日志和数仓各层明细数据HBase用来存实时特征结果和反欺诈名单支持毫秒级查询。计算层Spark负责离线批处理、特征工程、模型训练Hadoop MapReduce处理极端耗时的全量回溯任务。服务层规则引擎加载黑名单、硬性规则进行第一轮拦截评分模型对通过规则的申请输出信用评分封装成REST接口。展示层运营后台和大屏展示每日申请量、通过率、逾期率、评分分布等核心指标。整个架构的核心是中间那层特征宽表和评分模型。没有特征宽表Spark的算力无处安放没有评分模型前面的存储和计算只是存了一堆用不上的数据。这也是后面两章要重点展开的部分。2.3 集群部署起点伪分布式搭建、HA 架构和生产集群规划学习阶段不建议一上来就搞多节点集群。先在单机做Hadoop伪分布式搭建把HDFS的NameNode和DataNode、YARN的ResourceManager和NodeManager之间的关系跑明白再用Spark的local模式提交几个作业理解存储和计算的协作逻辑。集群里还有一个关键组件是ZookeeperHadoop和Zookeeper整合实战解决的核心问题是HA模式下NameNode的自动故障切换——主NameNode挂了备用节点要能自动顶上否则整个HDFS就瘫了。生产环境必须上HA架构这一点没有任何商量余地。NameNode是HDFS的单点元数据都在它内存里一挂全挂。Zookeeper负责协调两个NameNode的主备状态JournalNode负责同步编辑日志。生产集群的规模按数据量倒推1000万级注册用户每天新增日志约200GBKeep一个50个节点左右的集群存储和计算基本够用。节点规格建议是每台32GB内存、8核CPU、4块4TB硬盘这样的配置在大多数信贷业务里能从业务初期撑到中期。规模再大的话要考虑的是Spark任务的资源隔离而不是继续无限加节点。3. 数据接入与存储从原始JSON到可查询的特征宽表3.1 Spark中读取JSON日志的两种姿势和一个关键坑业务系统上报的日志绝大多数是JSON格式每条申请记录一个JSON文件或者一行一个JSON。Spark读取JSON最常见的错误是让Spark自己推断schema数据量大时这个推断过程会触发额外的扫描而且碰到包含多种字段形态的数据时推断结果经常和你预期不符——比如金额字段有时是字符串有时是数字Spark默认推断成string或bigint取出来才发现类型不对。我一般会直接在读取时显式指定schema代价是维护一个schema定义但换来的是稳定性和可预测性。典型的读取代码长这样from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, StringType, LongType, DoubleType, TimestampType spark SparkSession.builder \ .appName(credit_json_etl) \ .config(spark.sql.shuffle.partitions, 200) \ .enableHiveSupport() \ .getOrCreate() # 显式定义schema避免Spark自己推导导致的类型漂移 schema StructType([ StructField(user_id, StringType(), True), StructField(apply_time, TimestampType(), True), StructField(loan_amount, DoubleType(), True), StructField(device_id, StringType(), True), StructField(channel, StringType(), True), StructField(extra_info, StringType(), True) # 原始字段先保留后续解析 ]) # 读取HDFS上的原始JSON目录按天分区 raw_df spark.read \ .option(multiline, false) \ .schema(schema) \ .json(/data/raw/credit_apply/dt2024-01-01/) raw_df.printSchema() raw_df.show(5, truncateFalse)这段代码里multiline参数要按实际数据格式设置。每行一个JSON对象时设为false整个文件是一个大JSON数组时设为true。shuffle.partitions设置成200是一个比较保守的起步值它控制了Spark执行shuffle操作时的分区数量数据量大的时候这个值设置得太小单个任务处理的数据过多很容易撑爆executor内存。3.2 数仓分层ODS、DWD、ADS三层的建模思路原始数据直接拿来算特征是不现实的JSON里嵌套的字段解析要消耗大量计算资源而且重复解析会有一致性问题。标准做法是建立三层数仓ODS层原样保留原始JSONDWD层做清洗、脱敏、解析字段ADS层做特征宽表。ODS层就是上面读取的原始数据不做任何加工只按天分区。DWD层用Spark SQL做清洗和解析典型做法是用get_json_object把JSON里的嵌套字段提取出来同时对身份证号、手机号做脱敏处理。这一步很关键信贷数据涉及个人敏感信息审计时会要求你证明原始数据没有直接暴露在计算链路里。DWD层建表语句参考CREATE EXTERNAL TABLE IF NOT EXISTS dwd_credit_apply ( user_id STRING COMMENT 用户ID, apply_time TIMESTAMP COMMENT 申请时间, loan_amount DOUBLE COMMENT 申请金额, loan_term INT COMMENT 申请期限月, device_brand STRING COMMENT 设备品牌, channel_code STRING COMMENT 渠道编码, city_level INT COMMENT 城市等级 1-5, -- 脱敏后的手机号只保留前3后4 mobile_masked STRING COMMENT 脱敏手机号 ) PARTITIONED BY (dt STRING) STORED AS PARQUET LOCATION /warehouse/dwd/credit_apply; -- 按天覆盖写入分区 INSERT OVERWRITE TABLE dwd_credit_apply PARTITION (dt2024-01-01) SELECT user_id, apply_time, CAST(loan_amount AS DOUBLE), CAST(loan_term AS INT), get_json_object(extra_info, $.device_brand) AS device_brand, get_json_object(extra_info, $.channel_code) AS channel_code, get_json_object(extra_info, $.city_level) AS city_level, concat(substr(mobile, 1, 3), ****, substr(mobile, 8, 4)) AS mobile_masked FROM /data/raw/credit_apply/dt2024-01-01;建表时用STORED AS PARQUET的列式存储查询时只读取需要的列IO量可以降一半以上对特征计算阶段的多列读取特别友好。INSERT OVERWRITE的写法保证分区内数据是幂等的重跑任务不会产生重复数据。3.3 特征宽表为什么要宽以及怎么join才不翻车特征计算阶段最忌讳的是每次评分都去临时join十几张表。业界标准做法是预先算好一张特征宽表——每一行是一个用户每一列是一个特征。宽表的好处在于训练模型时直接把这张表喂给机器学习算法不需要再处理join逻辑实时评分时查询一次就能拿到一个用户的全部特征。构建宽表最常见的坑是join数据倾斜。比如用用户维度join授信记录某几个用户借款次数上千次会导致对应reduce任务数据量远大于其他任务表现为整体任务卡在99%。缓解手段有三个先把用户维度的数据做聚合压到每用户一条对join key做加盐处理分散热点或者干脆把低频维度做成广播变量。第三个手段在Spark里操作最简单核心代码如下from pyspark.sql import functions as F # 用户基本信息表数据量在几十万量级可以广播 user_info spark.table(dwd_user_info) user_features spark.table(dwd_credit_apply).groupBy(user_id).agg( F.count(user_id).alias(apply_cnt), F.avg(loan_amount).alias(avg_loan_amount), F.max(loan_amount).alias(max_loan_amount) ) # 广播小表避免shuffle user_features_with_info user_features.join( broadcast(user_info), onuser_id, howleft ) # 覆盖写入ADS特征宽表按天全量刷新 user_features_with_info.write \ .mode(overwrite) \ .format(parquet) \ .save(/warehouse/ads/user_feature_wide_table)这段代码里broadcast会把小表分发到每个executor的内存中join时不需要做shuffle这是Spark里最廉价的性能优化手段之一。宽表的刷新策略按业务容忍度来定信贷场景通常T1刷新就行——当天凌晨用前一天的数据重算全量宽表第二天白天做实时查询时读的就是最新版本。4. 特征工程与风险模型把Spark算力变成授信分4.1 特征计算与清洗Spark DataFrame的标准化数据处理特征宽表建好之后下一步是特征工程——把原始字段变成模型能用的数值特征。信贷风控常见的特征类型包括用户基本属性年龄、城市等级、职业类别、借款行为申请次数、平均借款金额、借贷间隔、设备信息设备使用时长、越狱/root标记、外部征信数据逾期次数、查询次数。这些特征不能直接进模型要做三类处理缺失值填充、异常值截断、数值标准化。Spark DataFrame做这套处理非常顺手代码示例如下from pyspark.sql import functions as F from pyspark.sql.functions import col # 读取特征宽表 feature_df spark.read.parquet(/warehouse/ads/user_feature_wide_table) # 填充缺失值数值列用中位数类别列用unknown stat feature_df.select( F.expr(percentile_approx(age, 0.5)).alias(age_median) ).collect()[0] feature_df feature_df.fillna({ age: stat[age_median], city_level: 3, # 缺失城市等级默认按三线城市处理 device_use_days: 0, channel_code: unknown }) # 异常值截断超过99分位数的值直接截断防止极端值拉偏模型 quantile_age feature_df.approxQuantile(age, [0.99], 0.01)[0] quantile_amount feature_df.approxQuantile(loan_amount, [0.99], 0.01)[0] feature_df feature_df.withColumn( age, F.least(col(age), F.lit(quantile_age)) ).withColumn( loan_amount, F.least(col(loan_amount), F.lit(quantile_amount)) ) # 标准化Z-score让不同量纲的特征处于同一尺度 mu_age, std_age feature_df.select( F.mean(age).alias(mu), F.stddev(age).alias(std) ).collect()[0] feature_df feature_df.withColumn( age_zscore, (col(age) - mu_age) / std_age )approxQuantile是Spark提供的近似分位数计算它不需要全量排序用抽样估算对于分位数这种统计量精度完全够。标准化用的是Z-score对于逻辑回归这类线性模型是必须的不标准化的话数值大的特征会在梯度计算中占主导地位。对于树模型标准化非必需但做了也没有副作用所以我一般统一做掉。4.2 规则引擎和评分卡模型怎么配合有了特征接下来是关键问题规则引擎和评分卡模型怎么分工。规则引擎处理那些一刀切的硬性风险——命中黑名单直接拒、申请金额超过限额直接拒、同一设备一天申请超过5次直接拒。规则引擎速度快、可解释性强、修改即时生效适合做第一道拦截。评分卡模型则处理那些介于中间的申请用一个分数来度量违约概率。评分卡最底层的模型通常是逻辑回归因为线性模型天然具备可解释性。在实际落地中为了更好处理非线性关系会先对连续特征做WOE分箱。这里用pyspark.ml.feature里的相关组件来实现from pyspark.ml.feature import QuantileDiscretizer # 对连续特征做分箱让模型学习非线性关系 discretizer QuantileDiscretizer( numBuckets10, inputColage, outputColage_bucket ) # 分箱结果转成OneHot编码后输入逻辑回归 from pyspark.ml.feature import OneHotEncoder ohe OneHotEncoder( inputCols[age_bucket, city_level_bucket], outputCols[age_ohe, city_ohe] ) # 组装特征向量并训练逻辑回归 from pyspark.ml.classification import LogisticRegression from pyspark.ml.feature import VectorAssembler assembler VectorAssembler( inputCols[age_ohe, city_ohe, apply_cnt, avg_loan_amount], outputColfeatures_vector ) lr LogisticRegression( featuresColfeatures_vector, labelColis_default, maxIter100, regParam0.01 )训练做完之后把模型的系数换算成分数每个分箱对应一个分数加总后映射到300-900分的信用分区间。分数越高违约概率越低。这个分数段业内默认是300分以下坚决拒、650分以上直接放、中间走人工复核。模型训练迭代的次数不用追求太多信贷场景里模型更新周期通常是季度级别特征是主导模型算法排第二。4.3 让评分接口算得快Spark内存参数与提交配置评分模型落地到线上服务常见的方式是把模型跑在Spark Streaming或者Structured Streaming上接Kafka里的申请事件流式计算实时给分。这时候性能参数的重要性不亚于模型本身的准确率。Spark Streaming消费Kafka数据做评分的典型配置参数如下from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(credit_scoring_streaming) \ .config(spark.executor.memory, 4g) \ .config(spark.executor.cores, 4) \ .config(spark.sql.shuffle.partitions, 400) \ .config(spark.streaming.kafka.maxRatePerPartition, 1000) \ .config(spark.sql.streaming.checkpointLocation, /warehouse/checkpoint/credit_scoring) \ .getOrCreate()spark.executor.memory决定了每个executor的JVM堆内存4g是生产环境的保守起步值别忘了堆外内存实际申请的物理内存要比这个值多20%左右。maxRatePerPartition限制每个Kafka分区每秒最多消费1000条防止流量突增直接把集群打爆。checkpointLocation一定要配流式任务的进度元数据全部存在这里不配的话任务重启后会丢数据或者重复消费。提交作业的方式也要配套。推荐用spark-submit提交到YARN集群而不是直接spark-submit跑在本地模式。提交命令参考spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 4g \ --executor-cores 4 \ --num-executors 30 \ --conf spark.sql.shuffle.partitions400 \ --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.2.0 \ credit_scoring_streaming.pynum-executors乘上之前配置的核数控制着整个作业的并行度。核数从4往上加收益会快速递减因为task调度和网络通信的开销也在增长。对大部分信贷数据量级来说30个executor已经能满足秒级评分的要求。如果评分延迟还是压不下来先看特征宽表有没有触发不必要的shuffle再考虑加资源这两者的顺序千万别搞反。5. 避坑清单Hadoop和Spark集群落地最常见的5个翻车现场5.1 数据倾斜导致某个executor内存溢出现象Spark作业卡在某个stage 99%不动日志刷出Container killed on request. Exit code is 143或者直接抛OOM。界面上看有一部分task执行时间异常长大部分task早就跑完了。原因特征宽表或者中间结果里某个key的数据量远大于其他key。信贷数据里最常见的倾斜源是user_id——部分高频用户或者渠道商账号对应的数据量是普通用户的几百倍这些数据全部落到同一个分区把对应的executor打爆。解决优先对倾斜key做过滤或单独处理比如单独处理借款次数超过100的用户如果倾斜程度可控给Spark配置加spark.sql.adaptive.enabledtrue和spark.sql.adaptive.skewJoin.enabledtrueSpark3.0以上的版本开启后会自动拆分倾斜分区最后的手段才是加内存。这一步的排查速度决定了你的加班时长我一般会把倾斜前后的数据量打出来对比五分钟定位。5.2 让Spark自己推断JSON的schema慢到怀疑人生现象读取一个几十GB的JSON目录Spark作业在读取阶段卡了很久日志显示一直在扫描文件。原因Spark为了推断出所有字段的类型需要先完整扫描一遍数据然后再真正读取一遍总共两遍IO。目录里文件数量多、嵌套层级深的时候时间成本非常可观。另一个副作用是推断结果不稳定数据形态稍微一变字段类型就跟着漂。解决第3章已经写过读取JSON时用.schema()显式指定。维护schema确实麻烦但比每次任务跑两遍要值得多。如果JSON里字段特别多可以先跑一次小样本推断出schema然后微调类型固化到代码里当常量。5.3 小文件过多把NameNode压垮现象集群没跑什么大任务但HDFS的NameNode频繁告警RPC响应延迟飙到几秒甚至几十秒。原因上游采集任务每个批次只写几十MB一天下来在HDFS上生成了上万个小文件。NameNode把每个文件、每个block的元数据都放在内存里文件数量一多内存占用上去读写请求的响应性能直线下降。Spark写特征宽表时如果分区设置过多且数据量不大也会产生同样的问题。解决一是上游把小块数据合并成大文件再上传控制单文件在128MB以上二是对已有的小文件跑一次合并任务。Spark侧的设置是控制输出分区数核心代码是这样# 写特征宽表前先重分区控制输出文件数量 from pyspark.sql import functions as F feature_df \ .repartition(50) \ .write \ .mode(overwrite) \ .format(parquet) \ .save(/warehouse/ads/user_feature_wide_table)repartition(50)把输出分片收缩到50个左右每个文件大约128MB这是一个平衡NameNode压力和后续读取并行度的折中值。文件太大读取并行度会下降太小元数据压力又上来。5.4 Zookeeper超时导致HA主备反复切换现象集群状态不稳定NameNode频繁自动切换甚至出现双主或同时无人服务的情况业务端隔几分钟就报一次HDFS连接失败。原因Hadoop和Zookeeper整合实战里最常见的一个坑——NameNode和Zookeeper之间的心跳超时设置过短。Zookeeper需要同时维护NameNode元数据、HBase协调、Kafka协调多个角色的会话任何一次GC暂停或者网络抖动超过超时阈值Zookeeper就会判定NameNode失联触发自动切换。切换本身又是高开销操作元数据加载需要时间于是形成切换-抢主-再切换的恶性循环。解决把zookeeper.session.timeout从默认的10秒左右调大到30秒dfs.namenode.avoid.read.stale.datanode加上让NameNode对DataNode的陈旧状态容忍度高一些。调参后的自检方法是杀掉active节点的进程观察standby节点能否在1分钟内接管然后恢复服务。这个动作要在业务低峰期做不然一个误判就直接影响线上审批链路。5.5 物理内存和容器内存对不上任务被莫名杀掉现象Spark任务运行到一半executor被YARN强制杀掉错误日志只有一行Container killed by YARN for exceeding memory limits。明明executor-memory配的不高怎么还被杀。原因Spark申请的内存默认只算JVM堆内内存但实际运行时还有堆外内存、线程栈、网络缓冲、Python进程如果用PySpark占用的内存。YARN限制的是整个容器进程的物理内存上限。申请4g的时候实际使用可能已经超过5g超过容器上限直接被kill。解决申请内存时预留20%-30%的buffer给堆外和Python进程比如需要4g的堆内就配spark.executor.memory4g的同时配spark.executor.memoryOverhead2g。YARN侧还要开启内存检测的宽松模式设置yarn.nodemanager.vmem-check-enabledfalse只检查物理内存不查虚拟内存。这套配置在很多博客里都语焉不详属于那种你以为配好了实际上没配的玄学参数。6. 拿到源代码和文档后怎么用先跑通最小闭环再改出你自己的风控系统拿到这套基于Hadoop、Spark的信贷风控系统源代码我建议你按三条线往下推不要上来就翻模型代码。第一条线是跑通数据链路——把生产环境的数据源换成你自己的样本数据哪怕是模拟的100万条先跟着文档把ODS到DWD的清洗任务跑完确认HDFS上能看到正确的分区数据。第二条线是跑通Spark作业——把特征宽表的生成脚本执行一遍看输出和文档里截图的数量级是否对得上这一步最好在伪分布式或单机测试集群上做。第三条线才是模型——用训练脚本跑一遍逻辑回归打开生成的模型报告确认覆盖率、区分度指标在合理范围。我最常踩的坑是把源码里写死的HDFS路径直接拿来用。文档里可能是/data/raw/credit_apply你的集群目录结构不一定一样而且建表的LOCATION权限、Hive的warehouse路径、Spark的checkpoint目录每一处都要按实际环境改一遍。改完之后用一个小技巧验证完整性跑一次数据抽样统计对比源文件的行数、字段数、空值率全部对上再继续下一步。至少先斩钉截铁跑通这个最小闭环你一晚上就能把整个系统的骨架摸清楚剩下的就是把规则阈值、模型参数按你的业务数据重新调一轮。希望帮到你。本文还有配套的精品资源点击获取