ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

信贷风控大数据实战:Hadoop+Spark全量特征工程与评分卡落地

信贷风控大数据实战:Hadoop+Spark全量特征工程与评分卡落地 简介一套基于Hadoop、Spark的大数据金融信贷风险防控系统的设计与实现资料包内含设计文档、完整源码及配套说明面向计算机相关专业学生、教师及企业开发者尤其适合毕业设计、课程设计或项目初期快速搭建演示环境。压缩包共70个文件以Java36个和Scala8个源码为主配合XML配置、SQL脚本及属性文件覆盖项目核心逻辑、资源配置与数据库初始化便于直接导入工程并运行调试。整体仅72KB轻量紧凑。项目已通过测试并获导师认可评审答辩得分95分可放心使用目前已有104人学习下载。借助该课题可掌握大数据平台与金融风控场景的结合方式在现有代码上二次开发快速完成课设或毕设方案。1. 信贷风控不是单机算法题HadoopSpark才能跑全量先放下“深度学习多厉害”的执念。信贷风控这个场景落地时最折磨人的往往不是模型而是数据根本跑不动。我见过有人用Pandas处理一亿条还款流水内存直接打满然后死机也见过有人抽样跑模型自认为AUC很漂亮上线后坏账率飙升。这份基于Hadoop、Spark的大数据金融信贷风险控系统资源就是冲这个痛点来的。它打包了完整设计文档、可运行源码、测试数据和部署脚本覆盖从Hadoop集群搭建到特征工程、评分卡建模再到阈值校准的全流程是典型的大数据课程设计和面试项目素材。它解决的是信贷业务里最实际的问题客户申请量上来了历史数据累积到千万甚至亿级怎么做全量特征计算、怎么训练可解释的评分卡模型、怎么把模型分数落到审批决策里。适合正在做大数据的课程设计、准备风控相关面试、或者刚接手公司风控数据平台需要参考落地方案的从业者。下面我把它拆成系统架构、特征工程、集群部署、踩坑记录、阈值校验五个维度代码部分可以直接抄作业。2. 架构分层与选型逻辑为什么离线全量计算必须挂Spark2.1 四层架构从数据采集到决策输出闭环怎么串起来先给这套系统画轮廓。整个架构划分成四层数据采集层、存储层、计算层、应用层。采集层对接进件系统、放款系统、第三方征信和前端埋点四类数据源日志用Flume落盘关系型库里的流水用Sqoop抽数交易事件走Kafka。存储层是HDFS加Hive数仓原始数据按业务日期做分区目录结构类似/user/hive/warehouse/ods/、dwb/、dws/这种数仓规范ods层存原始日志dwb层存清洗后的明细dws层存聚合后的指标。为什么不直接扔MySQL单日进件量上千万、特征上百维的时候MySQL一张宽表就要几百GB查询和扩缩容都成了瓶颈。计算层是整套系统最有分量的部分。ods层原始数据落盘之后要完成清洗、过滤、特征聚合、训练集拼接这几个步骤。老做法是MapReduce但MR每轮shuffle都要落盘特征工程做五六轮MR跑完一批亿级样本的特征宽表耗时能到两三个小时。换成Spark后DAG执行引擎把多轮依赖优化成一个图尽可能在内存里计算同一批任务压缩到20分钟上下。这套源码里做了两层计算协同Hive负责数仓ETLSpark负责特征工程和模型训练SparkSession初始化时通过hive.metastore.uris连通Hive Metastore直接读Hive表不用重复建表。应用层对应信贷风控的三个评分场景贷前申请评分A卡贷中行为评分B卡贷后催收评分C卡。源码把三种卡统一抽象成ScoringCard基类定义feature_engine、score、explain三个抽象方法A、B、C卡分别继承实现。这样做的好处是换数据集时只改特征工程部分评分流程和阈值判断逻辑完全复用。设计文档里有一张类图和一张时序图把一次在线审批从接口入参到查特征、算分数、返回决策的完整过程都画清楚了这部分对面试讲项目特别有用。2.2 选型逻辑Spark替代MapReduceHive与Spark共存Flink只做实时规则拆这套资源之前你可能会纠结到底该不该上Flink做实时风控我看完源码给出的判断是这套系统定位离线为主、实时规则为辅。信贷审批不是网络安全拦截一个进件申请从提交到出分秒级是合理的不需要毫秒级。而且评分卡的特征大量来自历史窗口聚合比如近六个月还款行为、近一年逾期次数天然是T1的离线计算。设计文档里能看同一个特征的离线批处理计算耗时统计——每日凌晨跑批做全量特征刷新白天业务查询直接读特征快照这个模式非常稳定。实时部分留了一个Spark Streaming模块处理两类不依赖模型的高置信规则单笔大额异常交易、设备指纹命中黑名单。实现方式是一个5秒的batch窗口从Kafka消费交易事件用SQL匹配规则命中后写入Redis在线接口秒查。但要提醒一句不要因为看到Spark Streaming就以为这套系统把模型打分搬到了实时链路。模型分数依赖大量窗口聚合特征实时逐条算成本太高正规做法就是离线算特征加在线查快照。这套源码的实时部分没有做模型推理是合理的取舍。选型还有一个关键点为什么保留Hive而不是纯Spark SQL替代因为数仓的元数据管理、分区管理、权限控制还是Hive Metastore这套生态成熟Spark SQL直接读取Hive表两者不是替代关系。用一张表把这几个引擎的边界说清楚对比维度MapReduceSpark RDDSpark SQLFlink中间结果落盘每轮都落盘尽量内存内存加自适应优化流式状态存储开发效率低代码繁琐中算子灵活高SQL表达能力强中API较陡典型场景超大文件简单清洗复杂迭代算法复杂join与聚合毫秒级实时风控在这套系统中的定位基本不用了少数UDF场景特征工程与模型训练未使用Streaming够用build_spark_session函数里有几个值得抄的配置spark.sql.shuffle.partitions设成600spark.sql.adaptive.enabled开启spark.serializer用Kryo。shuffle分区数设600是为了均衡并行度和单任务开销分区太少跑不满集群太多调度和序列化开销又把收益吃掉。在第5章避坑部分我会展开说这些参数到底踩了什么坑。2.3 源码模块划分接口抽象与SPI扩展看这个工程的组织方式能感觉到作者不是随便拼的。工程按功能分成risk-common、risk-etl、risk-ml、risk-streaming四个模块。common放日志、常量、配置工具类etl负责日常ETL作业ml放特征工程和模型训练streaming放实时规则引擎。模块之间通过接口依赖没有循环调用。有一点值得学习ml模块里的特征工厂用了SPI机制。新增一个特征时只需在META-INF/services下注册实现类不用改动既有聚合逻辑。这个扩展点在课程设计中不常见但放到真实风控团队里几乎每天用——业务要加一个新渠道特征你就加一个实现类跑批框架自动加载。文档里也明确写了“面向扩展开放面向修改关闭”。如果你是去面试能讲出这一点和只会背AUC的候选人差距立刻就出来了。对应到资源设计文档里面有数据字典、ER图、核心接口说明和部署拓扑。数据字典里给每个字段标注了来源表、口径定义和更新频率比如apply_cnt_30d的定义是“近30日客户提交贷款申请次数含被拒申请”。这里我想强调一下口径不一致是风控项目最容易吵架的地方。同样的“逾期次数”征信口径和业务口径可能差出一个量级设计文档把这些定义锁死了落地时不会扯皮。3. 特征工程与评分卡实现从JSON日志到A卡分数的完整链路3.1 用Spark SQL读JSON构建特征宽表半结构化日志怎么变成可观测量机器行为日志是最难处理的原始数据之一。它结构嵌套、字段不统一、量大。这套系统的做法是先spark.read.json把原始日志读入DataFrame再用select提取平铺。以下代码来自risk-etl模块的BehaviorFeatureJob我加了注释方便直接改改跑起来from pyspark.sql import SparkSession from pyspark.sql.functions import count, sum, avg, max, when spark SparkSession.builder \ .appName(behavior_feature_job) \ .enableHiveSupport() \ .config(spark.sql.warehouse.dir, /user/hive/warehouse) \ .config(spark.sql.shuffle.partitions, 600) \ .getOrCreate() # 按业务日期读取HDFS上的JSON行为日志分区目录就是日期 raw_df spark.read.json( hdfs://nameservice/warehouse/ods/behavior_log/2024-01-15 ) # 第一次接数据源时先打印schema确认字段别凭数据集文档写死列名 raw_df.printSchema() # 只保留分析需要的字段剔除埋点脏数据 behavior_df raw_df.select( customer_id, event_type, event_amount, event_time, channel ).filter( raw_df[customer_id].isNotNull() raw_df[event_type].isNotNull() ) # 核心聚合单一客户近30日的行为指标 feature_df behavior_df \ .filter(behavior_df[event_time] 2024-01-01) \ .groupBy(customer_id) \ .agg( count(when(behavior_df[event_type] loan_apply, 1)) .alias(apply_cnt_30d), sum(when(behavior_df[event_type] repayment, behavior_df[event_amount])) .alias(repay_amt_30d), avg(when(behavior_df[event_type] click, behavior_df[event_amount])) .alias(avg_click_amt_30d), max(event_time).alias(last_active_time) ) # 写入dws层宽表按客户id尾号做range分区减少小文件数 feature_df.write \ .mode(overwrite) \ .partitionBy(customer_id_tail) \ .bucketBy(32, customer_id) \ .saveAsTable(credit_db.dws_customer_feature)这段代码的逻辑拆开看。spark.read.json直接读HDFS目录目录一级就是日期分区跑批脚本按日期动态替换路径不需要在代码里做循环。printSchema这一步特别值钱JSON日志经常少了字段或改了类型不加这一步后面select会直接抛AnalysisException排错成本高得多。groupBy加agg的写法是典型的Spark SQL聚合模式。when条件里没匹配到的记录返回nullcount自动跳过null所以apply_cnt_30d统计的是“发生行为才计数”。event_time比较时直接用了字符串因为JSON日志里时间是ISO格式字典序就是时间序。参数上唯一要解释的是bucketBy(32, customer_id)——它按客户ID的哈希分到32个桶里下游和另一张同桶表join时可以走bucket join省掉一整轮shuffle。如果两张表桶数不一致Spark会退化成普通sort merge join收益就没了。3.2 WOE分箱与逻辑回归从模型系数到评分卡分数映射特征宽表就绪后进入模型训练环节。这里有个核心建模约定评分卡模型的输入是WOE值不是原始值。WOE全称Weight of Evidence衡量一箱内坏样本占比和单纯的数值归一化比WOE天然处理了缺失值和异常值且和逻辑回归线性输出衔接得很顺。分箱工具有现成的源码里用的是卡方分箱最小分箱数设4最大设8分完后手动把各箱的IV值打印出来做粗筛。from pyspark.ml.feature import VectorAssembler from pyspark.ml.classification import LogisticRegression from pyspark.ml.evaluation import BinaryClassificationEvaluator # WOE特征表已由分箱Job生成读Hive表 woe_df spark.sql( SELECT customer_id, is_bad, woe_age, woe_income, woe_debt_ratio, woe_apply_cnt_30d FROM credit_db.woe_features WHERE dt 2024-01-15 ) # 拼装特征向量输入列必须是WOE值 assembler VectorAssembler( inputCols[ woe_age, woe_income, woe_debt_ratio, woe_apply_cnt_30d ], outputColfeatures ) train_df assembler.transform(woe_df) # 逻辑回归训练regParam控制过拟合standardization对WOE值影响不大但建议保留 lr LogisticRegression( featuresColfeatures, labelColis_bad, regParam0.01, maxIter50, standardizationTrue ) model lr.fit(train_df) # 在训练集上看AUC只是确认模型收敛不能代表泛化能力 evaluator BinaryClassificationEvaluator( rawPredictionColrawPrediction, labelColis_bad, metricNameareaUnderROC ) train_auc evaluator.evaluate(model.transform(train_df)) print(f训练集AUC: {train_auc:.4f})逻辑回归训练过程的几个参数说明。regParam是L2正则项系数设0.01时对高维稀疏特征比较温和设大了模型偏向欠拟合设小了在强相关特征上容易出现系数爆炸。maxIter50对逻辑回归的BFGS优化器来说足够收敛不用调到几百。standardization设为True时Spark会在训练前把特征缩放为均值为0、方差为1WOE值本身是分布的压缩量标准化对系数影响不大但统一处理能防止某个特征量级过大干扰正则项。模型训练完之后还有一步逻辑回归系数到评分卡分数的映射。评分卡的标准形式是Score BaseScore Factor * ln(odds)其中odds是好人比坏人的比率。映射时先固定两个基准点在某个基础分数处odds为1比40每增加20分odds翻倍。源码里的ScoreCardCalculator类就是做这件事——把逻辑回归的截距项和每个变量的系数换算成每个分箱对应的分数值。这套公式在监管场景里很关键评分卡必须能解释成“偏离基准分多少分”黑盒输出一个神秘数字在审批场景是不被认可的。源码里还有个容易被忽略但很重要的设计特征穿越保护。它检查训练集的is_bad标签是否取了“未来”的数据做特征聚合。训练集切片时只用申请日期前的行为数据这个时间切分逻辑在DataSplitJob里是强制校验的。这块坑我在第5章细说。3.3 样本设计观察期、表现期与坏客户定义建模之前要先定义“谁是坏客户”这个问题定错后面全白干。这套资源的做法是观察期取申请日前12个月的行为数据表现期取申请日后6个月的还款表现。表现期内首次逾期超过30天就标记为坏客户is_bad1。以下是DataSplitJob和样本标注的简化版代码from pyspark.sql import SparkSession from pyspark.sql.functions import col, lit, when, datediff spark SparkSession.builder \ .appName(data_split_job) \ .enableHiveSupport() \ .getOrCreate() # 申请样本2023年6月到2023年12月的进件客户 apply_df spark.sql( SELECT customer_id, apply_date FROM credit_db.apply_info WHERE apply_date BETWEEN 2023-06-01 AND 2023-12-31 ) # 表现期内逾期超过30天的客户标记为坏样本 overdue_df spark.sql( SELECT customer_id, MIN(overdue_date) AS first_overdue_date, MAX(overdue_days) AS max_overdue_days FROM credit_db.repayment_schedule WHERE overdue_days 30 GROUP BY customer_id ) # 关联并生成训练标签 sample_df apply_df.join( overdue_df, customer_id, left ).select( apply_df[customer_id], apply_df[apply_date], when( col(first_overdue_date).isNotNull(), lit(1) ).otherwise(lit(0)).alias(is_bad) ) # 坏样本占比检查低于1%说明表现期短或坏客户定义过严 bad_rate sample_df.filter(col(is_bad) 1).count() / sample_df.count() print(f坏样本占比: {bad_rate:.4f})这里最关键的参数是overdue_days 30和表现期6个月。不同业务对“坏”的定义差别很大现金贷可能用逾期15天信用卡偏好用逾期60天。表现期长短直接影响样本量表现期设12个月样本时效性差设3个月很多坏客户还没暴露。这套源码里把表现期做成配置项跑批前在application.conf里改performance.period.months就行不用改代码。样本设计做完后训练集、验证集、测试集按时间切分不随机切。随机切分会让训练集和测试集混入同一时期的客户模型在时间维度上直接穿越AUC虚高得离谱。按时间切分才能模拟真实上线场景——用上个月的客户训练在这个月的新客上验证。这个细节也是面试里高频考点能主动说出来面试官会默认你真跑过项目。4. 集群搭建与作业提交从Hadoop伪分布式到Spark on YARN4.1 伪分布式搭建与环境变量第一步不要跳资源附带的部署文档从单机伪分布式起步。伪分布式模式让NameNode、DataNode、ResourceManager都跑在一台机器上这是学习和调试最快的方式也是很多人初次接触hadoop伪分布式搭建的地方。先看环境变量配置这一步错了后面全是报错# 配置JAVA_HOME、HADOOP_HOME和SPARK_HOME # 注意hadoop jar包版本必须和Spark编译版本匹配hadoop3对应spark的hadoop3版本 export JAVA_HOME/usr/lib/jvm/java-1.8.0-openjdk export HADOOP_HOME/opt/hadoop-3.3.4 export SPARK_HOME/opt/spark-3.3.0-bin-hadoop3 export HADOOP_CONF_DIR$HADOOP_HOME/etc/hadoop export PATH$PATH:$JAVA_HOME/bin:$HADOOP_HOME/bin:$HADOOP_HOME/sbin:$SPARK_HOME/bin # 验证环境 hadoop version spark-shell --version环境变量配置里最容易翻车的点JAVA_HOME必须指向JDK版本而不是只指向/usr/lib/jvm父目录。Hadoop 3.x要求JDK 8以上有些发行版默认JDK 11导致Hadoop开启参数不生效。HADOOP_CONF_DIR这个变量尤其重要Spark on YARN提交任务时要从这个目录读core-site.xml和hdfs-site.xml这个变量没配Spark就会用默认配置访问hdfs://localhost:9000而集群实际端口是9820连接直接拒绝。伪分布式的配置文件是core-site.xml和hdfs-site.xml。核心设置是把默认的HDFS地址指向本机然后数据块副本数设为1!-- core-site.xml -- configuration property namefs.defaultFS/name valuehdfs://localhost:9820/value /property property namehadoop.tmp.dir/name value/opt/hadoop/tmp/value /property /configuration!-- hdfs-site.xml -- configuration property namedfs.replication/name value1/value /property property namedfs.namenode.name.dir/name value/opt/hadoop/namenode/value /property property namedfs.datanode.data.dir/name value/opt/hadoop/datanode/value /property /configuration伪分布式里的dfs.replication1是关键。集群只有一台机器默认副本数3会导致DataNode上报块失败NameNode进入安全模式表现为hdfs dfs -ls卡死或报NameNode is in safe mode。hadoop.tmp.dir不显式配置的话Hadoop会用默认的/tmp目录系统重启后数据丢失这是另一个隐藏的坑。格式化文件系统后启动# 首次启动前必须格式化NameNode之后不要重复执行 hdfs namenode -format # 启动HDFS和YARN start-dfs.sh start-yarn.sh # 验证进程 jpsjps输出里必须能看见NameNode、DataNode、ResourceManager、NodeManager这几个进程。如果只有NameNode没有DataNode大概率是dfs.datanode.data.dir没有写权限或者格式化时集群ID不一致。hdfs namenode -format不要重复执行第二次格式化产生新的clusterIdDataNode里的旧clusterId对不上就会一直拿不到数据块。4.2 Hadoop HA与ZooKeeper整合NameNode高可用的关键配置从伪分布式往生产环境过渡第一件事就是把NameNode单点干掉。这台系统的部署文档里有一章专门讲hadoop和zookeeper整合实战。HA架构需要两台NameNode一台Active一台Standby共享状态通过JournalNode同步。ZooKeeper在这里负责两件事一是维护Active NameNode的锁二是触发自动故障切换。HA改造的几个核心配置项!-- hdfs-site.xml 中关于HA的部分 -- configuration property namedfs.nameservices/name valuemycluster/value /property property namedfs.ha.namenodes.mycluster/name valuenn1,nn2/value /property property namedfs.namenode.rpc-address.mycluster.nn1/name valuenode01:9820/value /property property namedfs.namenode.rpc-address.mycluster.nn2/name valuenode02:9820/value /property property namedfs.namenode.shared.edits.dir/name valueqjournal://node01:8485;node02:8485;node03:8485/mycluster/value /property property namedfs.client.failover.proxy.provider.mycluster/name valueorg.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider/value /property property namedfs.ha.automatic-failover.enabled/name valuetrue/value /property /configuration需要说明的参数。dfs.nameservices定义逻辑名称之后客户端连接hdfs://mycluster/而不是某个具体节点NameNode切换对业务透明。dfs.namenode.shared.edits.dir指向JournalNode集群两个NameNode通过这个共享存储同步编辑日志。dfs.ha.automatic-failover.enabled打开后ZooKeeper的ActiveStandbyElector会监控NameNode状态一旦Active节点心跳超时Standby节点自动切换。部署脚本里初始化HA要严格按顺序走# 1. 启动ZooKeeper集群 zkServer.sh start # 2. 在ZooKeeper中初始化HA状态 hdfs zkfc -formatZK # 3. 启动JournalNode、两个NameNode然后在Active节点格式化并初始化共享存储 hdfs namenode -initializeSharedEdits # 4. 启动备节点NameNode hdfs namenode -bootstrapStandby这里最容易错的是步骤2和3的顺序颠倒。先格式化ZK再初始化共享编辑日志顺序乱了会导致两个NameNode争抢active锁脑裂双写。ZKFC之间没有强制的锁协商顺序靠的就是初始化时序。我把这个坑踩过一次之后就养成了习惯每次HA部署都按文档顺序执行不凭直觉调整。4.3 Spark集群搭建与YARN提交参数作业参数决定了任务生死Spark集群搭建在配套的部署文档里占了很大篇幅。一个可用的Spark集群节点角色分成Master和Worker但生产上很少用Spark自带的Standalone模式大部分公司会把Spark跑在YARN上让YARN统一管理资源。这套源码的部署脚本也是这个思路Spark只作为计算框架接入YARN不启用Standalone。先看Spark on YARN的提交命令这是你每天都要打交道的# Spark on YARN提交示例 spark-submit \ --master yarn \ --deploy-mode cluster \ --class com.risk.ml.feature.BehaviorFeatureJob \ --name risk_feature_job \ --num-executors 12 \ --executor-cores 4 \ --executor-memory 16G \ --driver-memory 4G \ --conf spark.sql.shuffle.partitions600 \ --conf spark.sql.adaptive.enabledtrue \ --conf spark.serializerorg.apache.spark.serializer.KryoSerializer \ risk-ml-1.0.jar # 查看作业运行状态 yarn application -list yarn logs -applicationId application_1698888888888_0012提交参数的含义必须真正理解。num-executors12是executor数量executor-cores4是每个executor的CPU核数executor-memory16G是每个executor的堆内存。总资源等于三者相乘再加driver内存。关键约束是executor核数不要超过5超出后每个核心分到的内存太小垃圾回收频繁整个任务反而变慢。内存参数里有个隐含关系executor-memory设置的是堆内内存spark.memory.overhead会额外申请一块堆外内存默认是executor内存的10%但大shuffle场景10%不够。发生Container killed by YARN for exceeding memory limits时第一件事就是把spark.memory.overhead调大到2G或3G具体值看单个executor的数据膨胀程度。这个报错每年要吃掉很多新人一下午YARN判定的是物理内存超限不归JVM管。--deploy-mode cluster和client的区别也要搞清楚。client模式下driver跑在提交机器上适合调试日志直接打在控制台cluster模式下driver跑在YARN容器里适合生产调度日志要yarn logs去拉。课程设计阶段调试用client生产跑批用cluster不要混着用。5. 避坑排查数据倾斜、特征穿越与资源竞争的五个实战记录5.1 数据倾斜热点客户让某个Executor跑了40分钟现象同一个特征聚合Job大部分Executor十几分钟跑完有一个Executor卡了两个小时整个作业的进度条卡在99%。日志里出现ShuffleWriter写不出数据或者GC频繁。原因某个头部客户的行为数据量占了全量的30%。行为日志按customer_id做groupBy聚合时这个客户的所有事件都hash到同一个分区对应Executor要处理比平均大几十倍的数据量。风控场景里真存在这种客户高频小额借贷的忠实用户一年能产生几十万条事件。解决给特征聚合加两级处理。第一级先按customer_id聚合出粗粒度指标第二级单独处理超大客户。检测方式是在聚合前先按客户id count一下超过阈值量级就把该客户数据单独取出用repartition(客户id, 时间桶)切分成多个子任务再合并结果。Spark 3.0以后开启spark.sql.adaptive.skewJoin.enabled也能自动处理join场景的倾斜但groupBy场景还是要自己拆。5.2 内存溢出Executor直接OOM作业失败重启现象作业跑了一半yarn application状态从RUNNING变成FAILED日志里看到Java heap space或Container killed on request。有时还会伴随Lost executor连续出现。原因两个常见因素。一是spark.sql.shuffle.partitions设置过小比如保持默认200但单次shuffle输出的看到数据量达到200G每个分区1G数据单executor内存根本放不下。二是broadcast join的表超过了广播阈值Spark自动退化成sort merge join大表join大表时内存峰值飙升。解决按数据量算分区数。经验公式是目标每个shuffle分区处理的数据量控制在100M到200M之间。比如总shuffle输出20G600个分区每个分区差不多33M安全。spark.sql.autoBroadcastJoinThreshold默认10M超过阈值的维度表不要硬广播。最稳妥的办法是把大维表按客户id做bucket join两边桶数一致避免大shuffle。调整后如果还OOM再加spark.memory.overhead。5.3 特征穿越训练时AUC 0.85上线后只有0.62现象模型在训练集上AUC达到0.85KS接近0.55推上生产后实际区分能力大幅缩水坏账率不降反升。原因这是最典型的特征穿越也叫标签泄漏。训练标签用了申请日之后的数据做特征聚合。比如做A卡训练时特征宽表里混入了一个“近30天还款金额”字段但这个字段统计的区间覆盖了表现期。模型在训练时“偷看”了未来信息测试时同样因为时间窗口重叠而虚高一旦真正用于新申请这个未来特征根本不存在。解决DataSplitJob里强制做了时间边界校验。特征聚合只允许使用apply_date之前的数据标签只允许使用apply_date之后的表现期数据。跑批前自动检查特征表的event_time最大日期是否小于等于训练集最小申请日期超过就抛异常终止。落地习惯是每次训练前跑一遍校验脚本把特征表的日期上界和标签集的时间下界打印出来人工确认一遍。从那以后我每次训练前都强制走一遍这个校验再也不敢跳过。5.4 动态分区插入产生上千个小文件HDFS被文件数拖垮现象写Hive表时用动态分区跑完发现目标目录下有上千个几十KB的小文件后续查询全表扫描耗时成倍增加NameNode内存也被大量文件元数据占据。原因动态分区写数据时每个分区由一个或多个task各自写文件。shuffle分区数600目标表分区300个每个分区平均被两个task写每个task都产生了一个文件最后生成数百个。小文件问题不止影响查询也影响下一次写入的容错和恢复。解决写入前先repartition(分区数, 分区字段)让每个分区的数据尽量汇聚到少数几个task再写。再配合coalesce把写入文件数压到每个分区1到2个文件。实际参数是# 动态分区写入前按分区字段重新分区控制输出文件数 feature_df.repartition(64, customer_id_tail) \ .write \ .mode(overwrite) \ .partitionBy(customer_id_tail) \ .saveAsTable(credit_db.dws_customer_feature)repartition(64, ...)的64是按目标分区的数量估算的一个分区对应一到两个文件。分区多时这个值要上调分区少时强制降到和分区数接近避免生成大量空文件。写完以后用hdfs dfs -du -h检查目录下的平均文件大小低于128M的块大小就要重新审视分区和reduce的数量。5.5 Executor心跳丢失作业频繁失败日志却说网络正常现象大作业跑半小时后连续报几个Executor heartbeat timed out随之executor被强杀作业失败。检查网络和端口都是通的业务方也确认集群网络正常。原因这是YARN和Spark的资源竞争问题。executor的堆内存被GC长时间停顿占用时心跳线程无法及时发出心跳YARN判定executor失联。具体场景通常是某个executor在做大shuffle时GC停顿超过心跳超时阈值造成的。spark.network.timeout默认是120秒大作业shuffle时GC停顿可以超过这个值。解决调整心跳和网络超时参数同时给executor留出足够的堆外内存。常用做法是把spark.network.timeout增大到600秒spark.executor.heartbeatInterval保持默认或调大到30秒。另外确认spark.memory.overhead至少是executor内存的15%以上给netty的direct memory留余量。调完后还要观察YARN的ResourceManager日志确认没有节点被标记为blacklisted。6. 分数阈值校准用KS曲线找最佳切分点评分卡模型训练好AUC也到了0.78这时候还不能上线。评分卡输出的分数只是一个相对排序审批决策需要把它切成“通过”“人工”“拒绝”三个区间。源码里专门有一个ThresholdCalibrationJob用来在验证集上找最佳分数阈值。这一步做不好前面所有工程努力都白费。常规做法是跑一条KS曲线。KS统计量是每个分数档位上累计正样本率和累计负样本率的差差值最大的点就是区分能力最强的切分点。实现方式不复杂from pyspark.sql.functions import col, count, sum, when, row_number from pyspark.sql.window import Window # 读取验证集的模型预测结果 pred_df spark.table(credit_db.score_result) pred_df pred_df.withColumn( score_bucket, (col(score) / 20).cast(int) * 20 ) # 每个分数段统计好坏客户数 bucket_stats pred_df.groupBy(score_bucket).agg( count(when(col(is_bad) 1, 1)).alias(bad_cnt), count(when(col(is_bad) 0, 1)).alias(good_cnt) ).orderBy(score_bucket) # 计算累计占比并生成KS值 total_bad bucket_stats.selectExpr(sum(bad_cnt)).collect()[0][0] total_good bucket_stats.selectExpr(sum(good_cnt)).collect()[0][0] ks_df bucket_stats.withColumn( cum_bad_rate, sum(bad_cnt).over(Window.orderBy(score_bucket)) / total_bad ).withColumn( cum_good_rate, sum(good_cnt).over(Window.orderBy(score_bucket)) / total_good ).withColumn( ks, col(cum_bad_rate) - col(cum_good_rate) ) ks_df.select(score_bucket, ks).orderBy(ks, ascendingFalse).show(5)这段代码的评分逻辑是每个分数段的累计好坏差异最大化。分数段分桶宽度设20分一档先按桶汇总再用窗口函数算累计占比最后KS值最大化所在的分数段就是候选阈值。实际业务中还要结合两个约束通过率不能超过审批部门设定的上限拒绝率要把坏账率压到容忍范围内所以最终阈值往往不是KS最大的点而是“KS曲线进入平台期后的第一个点”。拿到候选阈值后回放历史申请数据模拟这个阈值下的通过率和预期不良率。用验证集回放是常规操作但要注意回放用的历史数据和上线后的客群分布可能偏移分群看PSI如果公共客群的分数分布变动超过0.1就需要重新校准。这套源码里把阈值校准做成一个独立Job每次模型迭代都要跑一遍状态存到HDFS的/result/threshold路径审批系统启动时读取这个配置。这套资源最打动我的一点就是它把容易被一笔带过的环节都做扎实了。设计文档讲清了架构权衡源码把特征穿越防护和阈值校准做成了强制步骤测试数据覆盖了常见的数据形态。我接手过的不少风控项目恰恰是倒在“模型好但工程糙”上。希望这份手记能帮你在复现时少走几个月的弯路也希望帮到你——照着设计文档跑通一遍全流程再回头看那些坑你就明白部署的每一步都在解决什么。本文还有配套的精品资源点击获取
返回列表