☰
信贷风控大数据实战:Hadoop与Spark构建特征工程和评分模型
2026/10/7 3:36:39 网站建设 项目流程

简介:一套基于Hadoop与Spark构建的金融信贷风险管控系统毕业设计实现方案,面向计算机专业正在开展毕业设计、课程设计或期末综合作业的本科生与学习者,聚焦信贷业务中的风险识别与管控问题。方案整合Hadoop的分布式存储与批处理能力、Spark的高效内存计算引擎,覆盖数据采集、清洗、特征提取、模型训练与多维度风险评估的全流程,并包含用户管理、数据接入、风险指标计算、可视化展示等核心功能模块,完整呈现大数据技术在金融风控领域的落地路径。压缩包共83个文件,以36个Java源码与8个Scala程序为主,配合12个XML配置、5个Properties环境配置,以及SQL脚本、前端JS与JSON辅助文件、README说明与工程备份,整体仅64KB,Maven工程结构清晰,便于按模块部署与二次开发。目前已有104人学习下载,代码经系统化调试与严格测试,组件接口衔接完整,运行环境配置清晰,既可为毕业设计或课程实践提供规范参考,也能帮助学习者快速理解Spark与Hadoop协同工作的工程实现。

1. 信贷风控上大数据:为什么 Hadoop 与 Spark 成了“标配”

做金融信贷风控的人迟早会撞上同一个坎:单机跑不动。客户申请记录、第三方征信查询、还款流水、APP 行为日志,一天少说几千万条,特征加工动辄要扫全量历史数据。早期靠 Oracle 存储过程跑批,一张宽表要跑四五个小时,模型上线后特征延迟一天半,等分数出来客户早去别家借完了。换分布式是唯一的出路,而 Hadoop 与 Spark 的组合是目前落地最顺的一条路:HDFS 负责把数据家底盘清楚,Hive 管数仓分层,Spark 做清洗、特征计算和模型训练。这套方案不是单纯的“上了大数据平台”,而是围绕信贷业务的生命周期——贷前审批、贷中监控、贷后催收——把数据加工和模型产出变成一条能稳定跑批的流水线。适合手里有数据但没过亿级别、团队刚从单机转分布式、想在 3 到 6 个月内把风控特征和评分模型搬到 Hadoop + Spark 上的数据工程师和风控建模人员。

2. 选型理由:信贷数据的“烂摊子”只有 Hadoop 能接,Spark 让建模不至于等死

2.1 为什么存储层要选 HDFS + Hive:半结构化数据与不可变日志的天然归属

信贷风控的数据源比互联网日志还要杂。除了核心系统导出的结构化申请表,还有第三方征信返回的 JSON 报文、合作渠道推送的 CSV 对账文件、App 端埋点上报的嵌套行为日志。这些数据归纳起来有一个共性问题:schema 不稳定。上游合作方这个月多加一个字段,下个月改了个枚举值,传统关系型数据库要么锁表改 DDL,要么先塞进一个 text 字段里等后面解析。HDFS 的好处是“先存后 schema”,文件落到 HDFS 上之后,Hive 建表时你想怎么解释都行,字段对不上最多查询出来是 NULL,不会把整个写入链路卡死。我们当时把来自 7 个渠道的进件数据原样丢进 HDFS,用 Hive 的贴源层保持原始字段名,后面分析时再层层裁剪,这种模式让数据接入周期从两周压缩到两天。

HDFS 的另一个优势是追加写和不可变性。信贷的申请记录、还款流水本质上都是事件日志,只增不改。传统数据库要应付 update、delete 带来的索引维护开销,在 HDFS 上根本没这个负担——新数据直接追加到新的数据块,NameNode 只管元数据,数据节点只管顺序写。用 Hive 做数仓分层时的典型做法是 ODS 层放原始数据,DWD 层做清洗过滤,DWS 层做轻度汇总,ADS 层给模型和报表直接查询。每一层用 insert overwrite 重建,跑挂了就从上一层重放,这种“后悔药”能力在单机数据库里是不可想象的。

HDFS 的容错也是信贷场景特别看重的。NameNode 做了 HA 之后,挂一个节点数据不丢;DataNode 的副本机制默认 3 份,机器磁盘坏了直接从副本恢复。信贷数据敏感,这 3 份副本不建议跨地域,但至少要在同一机房的不同机架,防止机架断电导致全部副本同时掉线。对于动辄存了 3 年全量历史数据的集群,这种容错能力直接决定了你敢不敢把历史数据全量保留。我们曾经为了省磁盘把副本数从 3 降到 2,结果一次机柜维护触发了大量数据块丢失,那周每天都在做 distcp 恢复数据,教训非常深刻。

2.2 为什么计算层要选 Spark 而不是裸写 MapReduce:迭代计算救不了建模

如果只是做 ETL,Hive on MR 也不是不行,但信贷风控逃不开迭代计算。特征工程阶段你要对全量历史数据反复扫描:先算客户的借款次数,再按月份切窗算均值,然后做一次去重和关联,最后再把特征拼成宽表。MapReduce 的每个 job 都要把中间结果落到磁盘,一个五步的特征计算流程要起五个 job,中间写五次 HDFS,一天的批处理任务能跑到深夜。Spark 的核心优势是把中间结果留在内存里,RDD 的血缘关系让多个计算步骤在一个 DAG 里完成,同样一套特征加工逻辑,Spark 跑 40 分钟,MapReduce 要跑 4 个小时,这不是夸张,是我们在同规模数据上实测的差距。

Spark 的另一个隐藏优势是 DataFrame API 和 SQL 的无缝切换。做风控特征的人大部分是写 SQL 出身,让他们直接写 Scala 或者 Java 的 MapReduce 是不现实的。Spark SQL 允许你用 SQL 查 DataFrame,也允许你在 SQL 里嵌套 UDF 做复杂逻辑,这意味着数仓团队已经写好的 Hive SQL 基本可以平移到 Spark 上跑,只是要执行引擎改成 Spark。我们的特征工程团队至今还在用 Spark SQL 写主要逻辑,只有涉及复杂窗口函数和自定义聚合时才会写 PySpark 代码,这个学习曲线对团队落地非常友好。

Spark on YARN 的资源调度也是选它的原因。信贷系统的批处理任务是典型的潮汐负载:白天业务系统压力大,批处理任务要避开高峰;凌晨到早上是黄金跑批窗口。Spark on YARN 可以和 Hive、MapReduce 任务混布在同一套集群里,YARN 根据队列配置动态分配资源,不需要为 Spark 单独搭一套集群。我们用的配置是生产队列和实验队列分开,生产队列跑正式批处理和模型打分,实验队列给建模人员跑探索性分析,两边互不挤占。资源不够时 Spark 会自己等待分配,不会像单机跑批那样直接把内存打爆。

2.3 系统整体架构:从数据接入到风险评分输出

整套系统的架构可以用六层来概括:数据接入层、存储层、计算层、特征层、模型层、服务层。数据接入层用 Flume 或 DataX 把业务库的数据同步到 HDFS,实时性要求不高的场景直接用 Sqoop 做每日全量或增量拉取;存储层就是 HDFS + Hive,按 ODS/DWD/DWS/ADS 四层组织;计算层跑 Spark,承担清洗、关联、特征计算的绝大多数任务;特征层把加工好的特征写回 Hive 表,同时导出一份到 Redis 里供线上决策引擎查询;模型层用 Spark MLlib 训练评分模型,产出模型文件后由打分程序加载,对每个客户算出一个分数;服务层是规则引擎和决策引擎,拿到分数后结合风控规则输出通过、拒绝、人工审核。

这套架构里最容易被忽略的是特征层和服务层之间的衔接。模型训练时用的特征和线上打分时用的特征必须完全一致,否则模型上线后性能必然衰减。我们的做法是特征加工逻辑只写一份,跑批生成特征宽表,线下训练直接查宽表;线上打分时用同样的跑批任务生成当日特征快照,推到 Redis 里供查询。跑批任务和线上特征任务共用同一套 Spark SQL,避免“线下一个版本、线上另一个版本”的经典翻车。

3. 数据接入与特征宽表构建:从 Hive 贴源表到 Spark SQL 特征加工

3.1 贴源层设计:Hive 分区表与字段规范

数据接入的第一步是把各个渠道的数据落到 Hive 贴源表里。信贷场景里最常见的分区策略是按天分区,分区字段叫dt,类型为 string,格式yyyyMMdd。按天分区的好处是跑批任务天然支持回溯——如果某一天的数据有问题,你把分区 drop 掉重新跑当天的任务就行,不影响其他分区。我们有些表还会加一个二级分区src_type,表示渠道来源,方便排查某个渠道单独的数据问题时直接定位到子分区做修复,不用全表扫描。

贴源表建表时有两个坑需要提前规避。第一个是字段类型不要跟着上游走,上游给你的金额字段可能是 string,但落到 Hive 里一定要转成 decimal——decimal(18, 2)是信贷金额的标准表达,避免浮点误差。第二个是原始报文不要急着解析,留一个raw_data字段存原始 JSON,后续解析有争议时还能翻出原文核对。这在和第三方征信数据对接时特别重要:征信报告的 JSON 嵌套四五层,你解析时难免有理解偏差,保留原始串相当于给自己留了后路。

CREATE TABLE IF NOT EXISTS ods_loan_apply ( apply_id STRING COMMENT '申请编号', cust_id STRING COMMENT '客户ID', product_type STRING COMMENT '产品类型', apply_amount DECIMAL(18,2) COMMENT '申请金额', loan_term INT COMMENT '借款期限(月)', apply_time STRING COMMENT '申请时间 yyyy-MM-dd HH:mm:ss', channel_source STRING COMMENT '渠道来源', raw_data STRING COMMENT '上游原始报文' ) PARTITIONED BY (dt STRING, src_type STRING) STORED AS ORC;

这段建表语句里,PARTITIONED BY是核心,dt 按天、src_type 按渠道,后续查询带上这两个分区条件能极大减少扫描量。STORED AS ORC是列式存储格式,信贷场景的查询大多是取部分列而不是整行,ORC 的列裁剪和压缩比默认的 TextFile 要好很多,我们实测同样数据量 ORC 比 TextFile 少占 70% 磁盘,查询速度也快一倍以上。如果你的集群还没装 ORC 的支持库,退一步用 Parquet 也可以,但不要用 TextFile 存贴源层,否则跑批任务会慢到你怀疑人生。

数据接入的调度框架,我们用的 Apache DolphinScheduler,工作流定义成“检查上游文件到达 → 执行 Sqoop 拉取 → 写入 Hive 分区 → 触发下游依赖”。核心原则是每个工作流节点都要带失败重试和告警,文件到达检查这一步尤其重要,上游渠道经常晚上 10 点才把文件推过来,你 8 点的调度任务等不到文件就挂了,没有重试机制的话当天数据就废了。

3.2 用 Spark SQL 做数据清洗:去重、格式统一与缺失值处理

贴源层数据落好之后,进入 DWD 层做清洗。常见清洗项包括:同一客户一天内的重复申请去重,日期字段格式统一成yyyy-MM-dd,金额字段负数和非数字字符串的处理,渠道枚举值映射成统一口径。用 Spark SQL 写清洗任务,本质是把一串可重复执行的 SQL 串到一个 Spark 作业里,比写一堆脚本再拼起来要直观得多。

INSERT OVERWRITE TABLE dwd_loan_apply PARTITION (dt = '${bizdate}') SELECT apply_id, cust_id, product_type, apply_amount, loan_term, FROM_UNIXTIME(UNIX_TIMESTAMP(apply_time, 'yyyy-MM-dd HH:mm:ss')) AS apply_time, CASE channel_source WHEN 'APP' THEN '01' WHEN 'H5' THEN '02' WHEN 'API' THEN '03' ELSE '99' END AS channel_source FROM ( SELECT *, ROW_NUMBER() OVER (PARTITION BY apply_id ORDER BY apply_time DESC) AS rn FROM ods_loan_apply WHERE dt = '${bizdate}' ) t WHERE rn = 1 AND apply_amount > 0;

这段 SQL 做了三件事:用ROW_NUMBER()窗口函数按申请编号去重,保留同一申请里时间最新的那一条;把时间字段格式统一;把渠道枚举值映射成两位编码。INSERT OVERWRITE保证当天分区重建,跑失败重跑也能保证结果一致。这里有一个容易踩的坑:去重字段如果上游不保证全局唯一,apply_id会有重复,去重逻辑必须写在 DWD 层,不要依赖上游保证。还有UNIX_TIMESTAMP解析失败会返回 NULL,如果有非标准时间格式的数据,整个 insert 任务会因为这个 NULL 关联产生脏数据——时间字段解析前可以先加一个WHERE apply_time RLIKE '^[0-9]{4}-[0-9]{2}-[0-9]{2}'过滤掉格式明显错误的行。

3.3 特征宽表加工:用窗口函数生成时序衍生变量

信贷风控的特征最看重“历史行为的时间切片”。比如客户近 3 个月的借款次数、近 6 个月最大逾期天数、近 1 年平均借款金额,这些特征是传统评分卡和机器学习模型都依赖的变量。特征加工的思路是先把业务数据按时间维度拆开,再用窗口函数在客户维度上聚合。

INSERT OVERWRITE TABLE dws_cust_loan_feature PARTITION (dt = '${bizdate}') SELECT cust_id, COUNTIF(loan_term >= 12 AND dt >= DATE_SUB('${bizdate}', 180)) AS cnt_term_ge12_last_6m, COUNTIF(dt >= DATE_SUB('${bizdate}', 90)) AS cnt_apply_last_3m, ROUND(AVG(apply_amount) OVER (PARTITION BY cust_id ORDER BY dt ROWS BETWEEN 180 PRECEDING AND CURRENT ROW), 2) AS avg_amount_last_6m, MAX(overdue_days) AS max_overdue_days FROM dwd_loan_apply WHERE dt >= DATE_SUB('${bizdate}', 365) GROUP BY cust_id;

这段 SQL 里的COUNTIF是 Spark SQL 3.0 之后支持的便捷写法,等价于SUM(CASE WHEN ... THEN 1 ELSE 0 END)。AVG(...) OVER (... ROWS BETWEEN ... PRECEDING AND CURRENT ROW)是一个滑窗函数,算的是截至当前日期的近 180 天平均申请金额,比先全量聚合再过滤更能反映时间衰减特性。实际做特征时要注意:滑窗的大小要和业务含义对齐——“近 3 个月”用 90 天,“近 6 个月”用 180 天,不要随手写个 100 这种没有业务含义的数字,否则后面模型解释性会很差。

特征宽表在 DWS 层做一次轻度汇总就够了,不要把每个特征的完整历史都存下来,否则数据膨胀得很快。更常见的做法是 DWS 层只保留“客户在某一天的特征快照”,后面建模时会拉到训练样本区间内每一天的快照做样本拼接。每天跑批的特征任务花 20 分钟,产出几千万客户的千维特征,这个规模在 Spark 上是非常轻松的。

4. Spark MLlib 训练信贷评分模型:坏客户定义、向量化与参数调优

4.1 坏样本定义与样本切分:表现期决定了模型的上限

信贷建模的第一步不是选算法,而是定义什么是“坏客户”。这个定义直接决定样本标签,进而决定模型学到的规律是否正确。常见口径是“逾期天数超过 30 天且逾期超过 90 天未还清”算坏客户,但不同产品线的容忍度不一样,现金贷的坏客户定义通常比大额分期更严。定义坏客户要配套一个概念:表现期。你在 2024 年 1 月放款的一批客户,需要至少观察 90 天才能判断他们是否变坏。因此建模时选的样本必须是“观察期结束且表现期已满”的客户,这个时间错位是信贷建模最容易翻车的地方——直接用最近几个月放款的客户做样本,标签大量缺失或未充分表现,训练出来的模型当然不准。

样本切分也要按时间切而不是随机切。随机切分会让训练集和测试集里的客户来自同一时间段,模型的时序泛化能力完全没被测试到。正确的做法是按放款月份切:比如用 2024 年 1 月至 6 月放款的客户做训练集,7 月至 9 月放款的客户做测试集,这样验证的是模型对未来新客户的表现。Spark MLlib 的randomSplit写起来爽,但在信贷场景一定不要用。

4.2 特征向量化与归一化:VectorAssembler 和 StandardScaler

Spark MLlib 的建模流程是 Pipeline 式的:先有一个 DataFrame 包含特征列和标签列,然后经过 Transformer 生成特征向量,最后喂给 Estimator 训练。特征列的类型必须统一成数值型,类别变量要先用StringIndexer+OneHotEncoder处理,连续变量要经过StandardScaler做标准化。信贷场景里各特征的量纲差距很大:申请金额可能到几十万,逾期天数最大只有几十,如果不做标准化,逻辑回归收敛会慢,基于距离的模型会被量纲主导。

from pyspark.ml.feature import VectorAssembler, StandardScaler from pyspark.ml.classification import LogisticRegression feature_cols = [ 'cnt_apply_last_3m', 'cnt_term_ge12_last_6m', 'avg_amount_last_6m', 'max_overdue_days', 'total_balance', 'monthly_income' ] assembler = VectorAssembler( inputCols=feature_cols, outputCol='features_raw' ) scaler = StandardScaler( inputCol='features_raw', outputCol='features', withStd=True, withMean=True ) lr = LogisticRegression( featuresCol='features', labelCol='is_bad', maxIter=100, regParam=0.01, elasticNetParam=0.5 ) pipeline = Pipeline(stages=[assembler, scaler, lr]) model = pipeline.fit(train_df)

这里StandardScaler的参数withStd=True表示除以标准差,withMean=True表示减去均值。对逻辑回归来说这两个都开效果最好;对树模型来说标准化不影响结果,因为树的切分只关心值的相对大小,省掉这一步也无妨。regParam是 L2 正则强度,elasticNetParam=0.5表示 L1 和 L2 各占一半,信贷特征维度高时弹性网既能稀疏化特征又能稳定解,比纯 L1 或纯 L2 更实用。训练完成后记得保存 pipeline 模型本身,而不是只保存 LR 系数,因为线上打分时要复用的是一整套特征处理流程,而不是一个裸模型。

如果特征列里有 NULL,VectorAssembler默认报错,直接让整个训练任务挂掉。信贷数据里特征缺失是常态,必须在组装向量之前统一处理:连续变量用中位数填充,类别变量用出现次数最多的值填充,或者干脆加一个is_missing标志列,让模型自己学缺失值的信息。我们实践下来,加缺失标志位比单纯填充效果更好,因为某些特征缺失本身就代表了客户的某种行为模式——比如收入字段缺失的客户,往往逾期率偏高。

4.3 模型评估与调参:用 AUC、KS 说话,别只看准确率

信贷数据集通常是严重的类别不平衡:好客户占比 95% 以上,坏客户可能只有 2% 到 3%。这种分布下准确率完全没有参考价值——你全猜“好客户”也能有 97% 的准确率。必须看 AUC 和 KS。AUC 衡量的是模型把好客户和坏客户分开的能力,0.5 等于瞎猜,0.7 以上才有业务可用性,0.8 以上算优秀。KS 统计的是好坏客户分数分布的最大差距,信贷行业惯例是 KS 大于 0.3 就可以试上线了。

from pyspark.ml.evaluation import BinaryClassificationEvaluator evaluator = BinaryClassificationEvaluator( labelCol='is_bad', metricName='areaUnderROC' ) train_auc = evaluator.evaluate(model.transform(train_df)) test_auc = evaluator.evaluate(model.transform(test_df)) print(f'Train AUC: {train_auc:.4f}, Test AUC: {test_auc:.4f}')

训练集和测试集 AUC 的差距如果超过 0.05,说明模型过拟合了。处理方法优先级从高到低:增加正则强度regParam、减少特征数量或做特征选择、增加更多训练样本。Spark MLlib 还提供了ParamGridBuilder做网格搜索,配合CrossValidator做交叉验证,但信贷场景里数据量大、模型参数空间小,手动调几组关键参数比全自动网格搜索更快更可控。我们通常只调三个参数:maxIter、regParam、elasticNetParam,对树模型则调maxDepth和maxBins。

如果换用随机森林,numTrees可以固定为 100 左右,再大收益很有限,训练时间线性增长。maxDepth从 10 开始试,太深容易过拟合,信贷特征大多没有那种需要极深切割的复杂非线性关系。maxBins决定连续特征离散化的桶数,默认 32 一般够用,但如果某个特征有超过 32 个离散取值且被当作数值列传入,信息会有损失——这种情况要单独处理而不是一味调大maxBins,因为桶数增加会显著拖慢训练速度。

5. 集群部署与跑批调度:一年踩过的 5 个坑和现在的稳定配置

5.1 伪分布式练手和真实集群的差距:Local 模式跑通不代表集群能跑

新手最常见的幻觉是 Spark Local 模式跑通了就万事大吉。热词里也总有人搜“伪分布式搭建”,因为教程多是从单机伪分布式开始教的。但伪分布式只验证了 API 和数据逻辑,完全不验证资源调度、数据倾斜、网络传输这些问题。常见做法是先在 3 台机器上搭最小集群,配置 YARN 和 HDFS,跑通一个简单的词频统计任务,再逐步加资源。我们团队新来的同事先在 Local 模式写了个特征计算,提交到集群后直接 OOM——本地是 8G 内存一个进程,集群上默认给 executor 申请 1G,内存不够当场失败。

从伪分布式切到真实集群要检查的配置有:spark.executor.memory和spark.executor.cores根据机器规格设置,spark.dynamicAllocation.enabled开启后让 Spark 根据负载自动申请和释放 executor,spark.shuffle.service.enabled在动态分配时必须开启,否则 executor 被回收后 shuffle 数据丢失。还有一个是spark.sql.shuffle.partitions,默认 200,对信贷这种几千亿行的大表关联,200 个分区会导致每个分区数据量巨大,shuffle 溢出到磁盘,任务慢成蜗牛。我们一般根据数据量调成 500 到 1000,并配合spark.sql.adaptive.enabled=true让 Spark 自动合并小分区。

5.2 NameNode 元数据膨胀:每天只跑一次备份不够

现象:集群运行半年后,NameNode 的 edit log 文件越来越大,偶尔出现“NameNode is in safe mode”告警,重启后恢复要等很长时间,期间整个 HDFS 不可写。

原因:HDFS 每次写操作都要记录 edit log,信贷集群每天的写入量是几千万个文件操作,edit log 持续膨胀。默认的dfs.namenode.checkpoint.period是 3600 秒,SecondaryNameNode 每小时做一次 checkpoint,但在写入高峰小时内的 checkpoint 可能来不及合并,log 会越积越大。

解决:把 checkpoint 周期从 3600 改成 600 秒,同时开启基于事务量的 checkpoint——dfs.namenode.checkpoint.txns设置为 100 万,哪个先触发就执行。同时把dfs.namenode.handler.count从默认 10 调高到 50,NameNode 在重命名、创建大量文件时 RPC 处理能力明显提升。我们就是这么改完之后,NameNode 进程稳定多了,满一年也没再进过 safe mode。

5.3 Spark 内存参数:executor 内存开再大也救不了数据倾斜

现象:跑特征关联任务时,某个 Spark Stage 一直卡在 99%,最后一个 Task 跑了 20 分钟还没结束,而其他 Task 几秒就完成了。

原因:数据倾斜。信贷数据里头部客户的记录数能占到总量的 20%——比如一个渠道放量时同一个渠道的申请记录几百万条,join 时按渠道分组,这个 key 所在的 reducer 就要处理几百万条数据,其他 reducer 只有几千条。调大 executor 内存只是把倾斜的后果推迟,不是根治。

解决:根治方案是加盐(salting)。对倾斜严重的 key 加一个随机前缀,把一条大 key 拆成多条小 key,join 完成后再去掉前缀。

from pyspark.sql.functions import concat, lit, rand, ceil # 左表加工:对倾斜 key 加盐 left_salted = left.df.withColumn( 'join_key_salted', concat(lit('t_'), (rand() * 10).cast('int').cast('string'), lit('_'), col('join_key')) ) # 右表放大:复制对应份数并加盐 right_expanded = right.df.withColumn( 'join_key_salted', concat(lit('t_'), explode(array([lit(i) for i in range(10)])).cast('string'), lit('_'), col('join_key')) ) result = left_salted.join(right_expanded, 'join_key_salted').drop('join_key_salted')

左表每条记录随机加一个 0 到 9 的前缀,右表对每条记录复制 10 份并分别加上 0 到 9 的前缀,这样原来的一个大 key 就分散到 10 个 task 上处理。这里rand()是 Spark SQL 的随机函数,explode把数组展开成多行。加盐份数不是越多越好,10 份通常够了,如果倾斜 key 的数据量特别大,调到 20 到 50 份也行,但会对右表产生 N 倍的存储和 shuffle 开销,要权衡。

5.4 小文件问题:Hive 分区表跑完一次批量插入,文件数多到 NameNode 压力大

现象:HDFS 上某张表的某个分区下有几千个小文件,每个只有几十 KB。跑查询时 Spark 启动几千个 Task 去读这些小文件,调度开销比计算本身还大,整个集群响应变慢。

原因:Spark 写 Hive 表时默认按分区写,每个 Task 生成一个文件。如果最后一步的 reducer 数量很多(比如spark.sql.shuffle.partitions设为 1000),就会生成上千个小文件。

解决:写 Hive 表之前增加一步合并小文件的 repartition,或者在写入后跑一次Hive的ALTER TABLE ... CONCATENATE。前者更可控,在 SQL 末尾加一句DISTRIBUTE BY dt, cust_id,让同一个分区和客户的数据落到同一个 reducer;后者对 ORC 格式的表最有效,是一条命令的事,不会重复数据。现在 Spark 3.x 的spark.sql.adaptive.coalescePartitions.enabled开启后能自动合并 shuffle 后的小分区,但只针对 shuffle 输出,对最后落盘文件的合并效果有限——我们最终是在调度工作流里加了一个定时任务,每天检查各表的小文件数量,超过阈值就自动执行CONCATENATE。

5.5 动态分区写入的坑:没有预先建分区,任务跑完分区表查不到数据

现象:Spark SQL 跑完插入任务日志显示成功,但查询WHERE dt = '20250201'结果为空,Hive 客户端SHOW PARTITIONS也看不到这个分区。

原因:动态分区写入前没有预先创建分区目录,Hive 的元数据没有记录新分区。Spark 往 Hive 表写数据时,如果表开了动态分区(hive.exec.dynamic.partition.mode=nonstrict),部分版本不会自动同步分区元数据到 Hive Metastore。

解决:写 SQL 之前先执行一条ALTER TABLE xxx ADD PARTITION (dt='20250201')把分区建好,再做INSERT OVERWRITE。或者更彻底的做法,写入时用spark.sql.catalogImplementation=hive搭配spark.sql.hive.metastore.partitionPruning之类参数,强制 Spark 同步元数据。我们后来干脆放弃了直接让 Spark 写 Hive 的动态分区表,改成INSERT OVERWRITE ... PARTITION (dt='${bizdate}')静态分区方式,稳定是第一位。

6. 让模型从“T+1 批跑”走向“分钟级”的风险决策:一个落地技巧

最后分享一个我们打磨了半年才稳定的实践:把模型打分从“每天早上出全量分”改成“增量客户实时算分”。信贷业务里有两类查询:一是贷前申请,客户刚刚点提交,你要在几秒内给一个评分;二是贷中存量客户,老客户在 App 上新增一笔借款,同样要快速判断能不能放款。全量批打分一天一次没办法覆盖这些场景,怎么办?我们的做法是:把 Spark 训练好的 pipeline 模型保存到共享存储,线上打分服务跑一个 Spark Streaming 任务,每 5 分钟消费一次新申请和存量客户的行为事件,在内存里加载模型,对增量客户做特征拼接和评分,结果写入 Redis,决策引擎直接查 Redis。

这个增量流水线有验证步骤要做:每周挑一天,把增量打分的结果和全量批跑的结果做一次全量比对,两条链路对同一批客户算出的分数差异必须小于阈值,否则说明特征时间窗口或模型版本有漂移。线上监控不看别的,只看分数分布和坏账率——用 PSI 看分数分布稳定性,用回看 AIC 看模型预测能力和实际结果的偏差。有一次我们发现线上分数整体抬高了 20 分,查了半天才发现是特征表里有一列近半年的逾期天数被上游改了数据来源,口径没对齐,从那以后我们给特征表加了字段血缘追踪,每次上游变更都要走一次全回归。做这套系统最大的感悟是:Spark 和 Hadoop 只是工具,真正决定风控系统价值的是特征口径的一致性和模型迭代的节奏感。希望帮到你。

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

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

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

立即咨询