
简介本资源是一套基于Hadoop与Spark架构的金融信贷风险管控系统源码面向计算机相关专业开展毕业设计、课程设计或项目实训的学习者帮助其在真实业务场景中理解分布式计算在金融风控领域的落地方法。压缩包共83个文件约64KB以36个Java源文件与8个Scala文件构成核心业务与计算逻辑辅以12个XML配置、5个properties参数文件、1个SQL建表脚本及README说明文档整体结构清晰、便于按模块研读。系统覆盖用户管理、数据接入、风险指标计算与可视化展示等模块实现从数据采集、清洗到特征提取、模型训练的全流程处理Hadoop负责海量数据存储与批处理Spark提供内存计算能力两者协同完成多维度信贷风险评估。目前已有104人学习关注适合需要规范项目参考与可扩展开发基础的学习者借鉴。1. 金融信贷风险分析为什么非要上 Hadoop 与 Spark一笔贷款背后的数据账一笔消费贷从申请到放款背后要过的数据关卡远比多数人想象的多。身份核验、征信查询、多头借贷扫描、设备指纹、行为埋点、历史还款表现、关联企业图谱——单笔申请跑完这些维度落下来的原始记录轻松过万条。一家中型城商行日均有几万笔进件一年就是几十亿条明细。这些数据放在单机 MySQL 里跑一次全量逾期率回溯要几个小时特征工程窗口一拉长直接 OOM。金融信贷风险大数据分析系统要解决的核心矛盾就在这里数据量级和计算复杂度同时上来了传统关系型数据库的存储和算力都撑不住。Hadoop 负责把海量明细沉到 HDFS 上做低成本存储Spark 负责在内存里做迭代计算和特征加工两者配合是目前业内最成熟、落地案例最多的一条路。这套方案适合正在做信贷风控平台选型的数据工程师、需要把离线评分卡跑通的后端开发以及要交大数据课程设计的学生。下面从环境搭建一路讲到调参和踩坑能照着复现。2. 从零搭一套能跑信贷数据的 Hadoop 与 Spark 环境伪分布式够用吗2.1 伪分布式与全分布式的选型判断先回答一个被问烂的问题做信贷风险分析到底要不要搭全分布式集群。我的判断标准很简单——看你的数据量和迭代频率。如果只是跑课程设计、验证评分卡逻辑、做几十万条样本的特征工程伪分布式完全够用一台 16G 内存的机器就能把 HDFS、YARN、Spark 全跑起来。伪分布式的本质是每个 Hadoop 守护进程单独起一个 JVMNameNode、DataNode、ResourceManager、NodeManager 都在本机数据副本数默认设成 1。它的好处是配置简单、排错直观坏处是没有真正的并行度Spark 的 executor 只能在本机抢内存。真正上生产、要处理千万级以上信贷明细、还要跑多轮模型迭代的必须上全分布式。三台起步一台 Master 跑 NameNode 和 ResourceManager两台 Worker 跑 DataNode 和 NodeManager。内存分配上NameNode 给 4G每个 NodeManager 按机器总内存的 70% 配 yarn.nodemanager.resource.memory-mb剩下的留给系统。这里有个血泪经验很多人把 NodeManager 内存配到 90%结果 Spark executor 一申请大内存容器系统直接触发 OOM Killer 把进程杀了日志里只看到容器退出码 137排查半天以为是代码问题。选型上还有一条如果团队已经在用云厂商的托管 Hadoop 服务别自己从零搭。托管服务帮你处理了元数据备份、节点扩容、安全认证这些脏活你只需要专注写 Spark 作业。自建集群的运维成本在节点超过十台后会指数上升。2.2 伪分布式环境的最小搭建步骤下面这套步骤在 Ubuntu 20.04 上验证过JDK 用 8 或 11 都行Hadoop 用 3.x 系列。先建一个专用用户别用 root 跑否则后面权限问题会让你怀疑人生。# 创建 hadoop 用户并配置免密 sudo useradd -m -s /bin/bash hadoop sudo passwd hadoop su - hadoop ssh-keygen -t rsa -P -f ~/.ssh/id_rsa cat ~/.ssh/id_rsa.pub ~/.ssh/authorized_keys chmod 600 ~/.ssh/authorized_keys # 下载并解压 Hadoop版本按官网当前稳定版选 wget https://downloads.apache.org/hadoop/common/hadoop-3.3.6/hadoop-3.3.6.tar.gz tar -xzvf hadoop-3.3.6.tar.gz -C /home/hadoop/ mv /home/hadoop/hadoop-3.3.6 /home/hadoop/hadoop解压完配置四个文件。core-site.xml 指定 HDFS 的默认文件系统地址hdfs-site.xml 设副本数为 1mapred-site.xml 指定用 YARN 跑 MapReduceyarn-site.xml 配 ResourceManager 地址。这几个文件的具体内容网上教程很多核心是保证 fs.defaultFS 指向 hdfs://localhost:9000yarn.nodemanager.aux-services 设成 mapreduce_shuffle。# 格式化 NameNode只做一次重复格式化会丢元数据 hdfs namenode -format # 启动 HDFS 和 YARN start-dfs.sh start-yarn.sh # 验证进程 jps # 应该看到 NameNode、DataNode、ResourceManager、NodeManager、SecondaryNameNodejps 输出里如果少了 DataNode九成是多次 format 导致 clusterID 不一致。解决办法是把 dfs/data 和 dfs/name 目录清掉重新 format别硬扛。2.3 Spark 的部署模式与信贷场景的匹配Spark 装的时候注意版本要和 Hadoop 对齐。下载页面上会标 pre-built for Apache Hadoop 3.3 and later 这种字样选对应的包省得自己编译。解压后配 spark-env.sh设 JAVA_HOME 和 HADOOP_CONF_DIR后者指向 Hadoop 的 etc/hadoop 目录这样 Spark 才能认出 HDFS 和 YARN。信贷风险分析场景下Spark 的部署模式选 YARN 而不是 Standalone。原因有两个一是 YARN 能统一管理 Hadoop 和 Spark 的资源不会出现两套调度器抢内存的情况二是 YARN 的队列机制可以给风控作业单独划资源池避免和离线报表任务互相挤。提交作业时用--master yarn --deploy-mode clusterdriver 跑在集群里客户端断了作业也不受影响。# 提交一个 Spark 作业到 YARN 的典型命令 spark-submit \ --master yarn \ --deploy-mode cluster \ --executor-memory 4g \ --executor-cores 2 \ --num-executors 10 \ --conf spark.sql.shuffle.partitions200 \ --conf spark.yarn.executor.memoryOverhead1g \ credit_risk_feature.py参数说明executor-memory 是每个 executor 的堆内存memoryOverhead 是堆外内存处理信贷数据里大量的字符串字段时堆外内存不够会报 Container killed by YARN for exceeding memory limits。num-executors 乘以 executor-cores 就是总并发核数别超过集群可用核数的 80%。shuffle.partitions 默认 200数据量小的时候反而拖慢可以按数据量除以 128MB 来估算。3. 信贷明细从 HDFS 到特征宽表Spark SQL 与 DataFrame 的加工链路3.1 原始数据的落库结构与分区策略信贷风险分析的原始数据一般分几类申请进件表、还款表现表、征信查询表、设备行为表。这些数据从业务库同步到 HDFS 时最常见的格式是 Parquet 或 ORC列式存储压缩比高Spark 读的时候能只扫需要的列。分区字段选日期按 dt2024-01-01 这种目录结构组织。分区的好处是回溯某个月的逾期率时Spark 只扫对应目录不用全表扫。有个容易翻车的点分区字段的基数不能太高。有人拿用户 ID 做分区几百万个目录NameNode 直接爆。分区字段的基数控制在几千以内日期、产品线、渠道这类维度合适。# 读取 HDFS 上的信贷申请明细按日期分区 from pyspark.sql import SparkSession from pyspark.sql.functions import col, datediff, when, count, sum spark SparkSession.builder \ .appName(credit_risk_feature) \ .config(spark.sql.parquet.compression.codec, snappy) \ .enableHiveSupport() \ .getOrCreate() # 读取申请进件表只取需要的列减少 IO apply_df spark.read.parquet(/warehouse/credit/apply/dt2024-01-*) \ .select(apply_id, cust_id, apply_time, loan_amount, loan_term, product_code, channel_code) # 读取还款表现表标记逾期 repay_df spark.read.parquet(/warehouse/credit/repay/dt2024-01-*) \ .select(apply_id, due_date, repay_date, repay_amount) # 计算每笔申请的逾期天数 overdue_df repay_df.withColumn( overdue_days, when(col(repay_date).isNull(), 999) .otherwise(datediff(col(repay_date), col(due_date))) ).groupBy(apply_id).agg( sum(when(col(overdue_days) 0, 1).otherwise(0)).alias(overdue_cnt), count(apply_id).alias(total_periods) )这段代码的逻辑是先把申请和还款两张表读进来只选分析需要的列。Parquet 的列裁剪在这里很关键信贷表动辄上百列全读进来内存直接翻倍。逾期天数的计算用 datediff还款日期为空说明还没还标记成 999 天。最后按申请 ID 聚合出逾期期数和总期数这两个字段是后面算逾期率的基础。参数上注意 spark.sql.parquet.compression.codec 设成 snappy压缩和解压速度快适合信贷这种需要频繁读取的场景。如果磁盘紧张可以换 gzip但读取会慢一些。3.2 特征宽表的构建与时间窗口处理风控特征的核心是时间窗口。近 3 个月申请次数、近 6 个月最大逾期天数、近 12 个月多头借贷平台数——这些特征都要按时间窗口算。Spark 里做窗口聚合用 Window 函数或者直接按时间范围过滤后 groupBy。前者适合窗口固定、数据量大的场景后者适合窗口灵活、需要频繁调整的探索阶段。from pyspark.sql.window import Window from pyspark.sql.functions import row_number, lag, avg # 按客户 ID 分区申请时间排序算历史申请次数 window_spec Window.partitionBy(cust_id).orderBy(apply_time) \ .rowsBetween(Window.unboundedPreceding, Window.currentRow) apply_df apply_df.withColumn(hist_apply_cnt, count(apply_id).over(window_spec)) # 算近 6 个月的平均申请金额 window_6m Window.partitionBy(cust_id).orderBy(apply_time) \ .rangeBetween(-180 * 86400, 0) # 按秒算的 180 天 apply_df apply_df.withColumn(avg_amount_6m, avg(loan_amount).over(window_6m))rangeBetween 的参数是秒数180 天乘以 86400 秒。这里有个坑apply_time 必须是 timestamp 类型如果是字符串要先 to_timestamp 转换否则 rangeBetween 会报类型错误。另外 rowsBetween 和 rangeBetween 的区别要搞清楚前者按行数算后者按值算时间窗口必须用 rangeBetween。特征算完落到 HDFS 上格式还是 Parquet按客户 ID 和申请日期分区。这张宽表就是后面跑评分卡和模型训练的输入。落表的时候注意字段类型金额用 decimal 别用 double避免精度丢失导致对账对不上。3.3 用 Spark SQL 做逾期率的多维分析特征宽表建好后业务方最常要的是各种维度的逾期率。产品线维度、渠道维度、地区维度、授信额度区间维度——这些用 Spark SQL 写起来比 DataFrame API 更直观。-- 按产品线和渠道统计逾期率 SELECT product_code, channel_code, COUNT(DISTINCT apply_id) AS apply_cnt, SUM(CASE WHEN overdue_cnt 0 THEN 1 ELSE 0 END) AS overdue_cnt, ROUND(SUM(CASE WHEN overdue_cnt 0 THEN 1 ELSE 0 END) * 100.0 / COUNT(DISTINCT apply_id), 2) AS overdue_rate FROM credit_feature_wide WHERE dt BETWEEN 2024-01-01 AND 2024-06-30 GROUP BY product_code, channel_code HAVING COUNT(DISTINCT apply_id) 100 ORDER BY overdue_rate DESC;HAVING 条件过滤掉样本量太小的组合否则某个渠道只有 5 笔申请、逾期 1 笔逾期率 20%这种数字没有统计意义反而误导业务判断。样本量阈值一般设 100 到 500看具体业务的数据密度。这段 SQL 跑在 Spark SQL 引擎上底层会转成 Catalyst 执行计划做谓词下推和列裁剪。如果发现跑得慢先看执行计划里有没有全表扫描再检查 shuffle 分区数是不是太少导致数据倾斜。某个产品线的数据量特别大时可以单独给这个维度加盐打散。4. 信贷风控作业的避坑与排查那些让任务跑挂的细节4.1 数据倾斜导致某个 task 卡死现象Spark 作业跑到 99% 卡住不动打开 Spark UI 看某个 task 的 shuffle read 是其他 task 的几十倍持续时间也长得多。信贷数据里这种情况常见于按客户 ID 聚合时某个大客户的申请记录特别多或者某个渠道的数据量远超其他渠道。原因shuffle 阶段按 key 分区同一个 key 的数据全落到一个 task 上。客户 ID 做 key 时如果某个客户有几万条记录这个 task 就要处理几万条其他 task 可能只有几十条。解决给倾斜的 key 加随机前缀打散。比如客户 ID 后面拼一个 0 到 9 的随机数先做一次局部聚合再去掉前缀做全局聚合。或者用 Spark 3 的 AQEAdaptive Query Execution开spark.sql.adaptive.enabledtrue和spark.sql.adaptive.skewJoin.enabledtrue让 Spark 自动处理倾斜。AQE 在信贷场景下效果明显建议默认开启。4.2 内存溢出与堆外内存不足现象任务报 java.lang.OutOfMemoryError: Java heap space或者 Container killed by YARN for exceeding memory limits。前者是堆内存不够后者是堆外内存不够。原因信贷特征工程里经常要缓存宽表做多轮迭代cache 的数据量超过 executor 内存。另外 Spark 的 shuffle 和 join 操作也会占用大量堆外内存。解决先看 Spark UI 的 Storage 页面确认缓存占了多少内存。如果确实需要缓存大表提高 executor-memory同时把 memoryOverhead 设成 executor-memory 的 10% 到 20%。如果还是不够考虑用persist(StorageLevel.DISK_ONLY)把缓存落到磁盘牺牲速度换稳定。另外检查代码里有没有 collect() 把大结果集拉到 driver这个操作在信贷场景下几乎必炸改成 write 落 HDFS。4.3 小文件过多拖垮 NameNode现象HDFS 上某个目录下有几十万个小文件每个几 KBNameNode 内存告警Spark 读的时候 task 数爆炸调度开销比计算还大。原因Spark 作业的并行度太高每个 task 写一个文件。或者上游数据同步时按小时切分一天 24 个文件一年下来小一万个。解决写完数据后用coalesce或repartition合并文件目标文件大小 128MB 到 256MB。已经产生的小文件用 Hadoop 的getmerge或者写个定时任务合并。建表时开 Hive 的合并参数hive.merge.mapfilestrue和hive.merge.mapredfilestrue让 Hive 自动合并小文件。4.4 时间窗口计算的结果对不上现象用 rangeBetween 算近 6 个月申请次数和业务方用 SQL 在 MySQL 里算的结果差了几十条。原因时间窗口的边界处理不一致。Spark 的 rangeBetween 是闭区间MySQL 的 BETWEEN 也是闭区间但时区可能不同。Spark 默认用 UTCMySQL 用本地时区跨天的时候差 8 小时窗口边界就错位了。解决统一时区。Spark 里设spark.sql.session.timeZoneAsia/Shanghai数据里的时间字段统一转成这个时区。另外确认 apply_time 的精度是精确到秒还是毫秒rangeBetween 的参数单位要和精度匹配。4.5 YARN 队列资源不足导致作业排队现象spark-submit 提交后一直处于 ACCEPTED 状态不进入 RUNNINGYARN 的队列资源被其他作业占满。原因风控作业和离线报表作业共用一个队列报表作业跑全量的时候把资源吃光了。解决给风控作业单独建 YARN 队列配最小资源保证。在 capacity-scheduler.xml 里配队列的 capacity 和 maximum-capacity风控队列的 capacity 设 30%maximum 设 50%保证报表跑满时风控还能拿到资源。提交作业时用--queue risk指定队列。5. 让信贷风险分析跑得更快几个我反复用到的调优习惯调优这件事我的习惯是先看 Spark UI 再动手别凭感觉改参数。UI 里三个地方必看Stage 页面的 task 时间分布如果 max 和 median 差几倍说明有倾斜SQL 页面的执行计划看有没有 CartesianProduct 或者全表扫描Storage 页面的缓存占用确认内存花在哪。第一个习惯是控制 shuffle 的数据量。信贷特征工程里 join 操作特别多申请表 join 还款表、join 征信表、join 设备表。每次 join 都是一次 shuffle。能提前过滤的数据先过滤能 broadcast 的小表用 broadcast join。征信表如果只有几万行直接broadcast(credit_df)避免 shuffle。Spark 3 里 broadcast 阈值默认 10MB小表超过这个数可以调spark.sql.autoBroadcastJoinThreshold但别调太大driver 内存扛不住。第二个习惯是分区数按数据量算。shuffle.partitions 默认 200数据量 10GB 的时候每个分区 50MB偏小task 调度开销大。按 128MB 一个分区算10GB 设 80 个分区就够。数据量 100GB 就设 800。分区数太少会倾斜太多会调度慢128MB 是个比较稳的中间值。第三个习惯是序列化用 Kryo。Spark 默认的 Java 序列化又慢又占空间换成 Kryo 后 shuffle 数据量能少 30% 左右。配置方式是spark.serializerorg.apache.spark.serializer.KryoSerializer然后把自定义的类注册进去。信贷场景里常用的 POJO 都注册上效果更明显。第四个习惯是监控 executor 的 GC 时间。Spark UI 的 Executors 页面能看到每个 executor 的 GC 耗时如果超过总时间的 10%说明内存压力大要么加内存要么减少缓存。信贷特征宽表字段多序列化后的对象也不小GC 频繁是常态定期看这个指标能提前发现内存问题。最后一个习惯是给关键作业加 checkpoint。信贷评分卡训练要迭代几十轮每轮都从 HDFS 重读数据太慢。在迭代开始前把特征宽表 cache 住中间结果定期 checkpoint 到 HDFS断了能续跑。checkpoint 的目录设在 HDFS 上别设本地否则 executor 挂了 checkpoint 就丢了。这些习惯不是一次配好的是每次作业跑挂之后加一条慢慢攒出来的。我现在接手一个新集群第一件事就是把 AQE、Kryo、时区这三个配置写进 spark-defaults.conf能省掉后面很多排查时间。希望帮到你。本文还有配套的精品资源点击获取