ARTICLE DETAIL

资讯详情

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

基于Hadoop+Spark的大数据信贷风控系统毕设全解析

基于Hadoop+Spark的大数据信贷风控系统毕设全解析 简介本资源是一套面向计算机及相关专业本科生的毕业设计级大数据风控系统实现聚焦金融信贷场景下的风险识别与评估问题适用于毕设、课程设计、项目实训及大数据入门实践。系统基于Hadoop生态构建数据存储与调度底座采用Spark进行分布式特征工程、模型训练与实时评分计算技术栈覆盖Scala开发、Maven工程管理、Spark Streaming流式处理及配置化数据源接入兼顾批流一体架构设计思想。压缩包共23个文件含8个核心Scala业务逻辑与作业类、7个XML配置文件含pom.xml与IDEA项目配置、2个properties环境参数文件以及README.md说明文档和前端交互相关JS/JSON文件整体仅53KB轻量易读。已有2377人学习下载代码经实测可运行提供完整目录结构与模块划分如data-source-spark-streaming子模块便于理解大数据风控系统的分层设计与工程落地路径。 每年毕业季都能看到一批同学把“XX管理系统”当毕业设计题目交上来的东西十有八九是Spring Boot加MySQL的增删改查。但同样是做系统“基于HadoopSpark的大数据金融信贷风险控系统”完全是另一个物种它要解决的不是“数据怎么存”而是“数据量大到单机顶不住、计算逻辑复杂到一张表装不下”之后整套信贷风控流程该怎么建模、怎么落代码、怎么向评委证明它真的在跑大数据。这个题目我前后带过几届学生源码结构、技术选型、答辩被追问的点基本都有规律可循。这篇文章就把这套系统从业务到代码到答辩的逻辑完整拆一遍适合正在做这个方向的毕设、或者想入门大数据开发但缺一个完整项目经验的同学参考。1. 这个题目真正在考什么从“信贷风控”到“大数据架构”的完整映射1.1 信贷风控的本质在放款之前就想清楚钱能不能收回来很多同学拿到这个题目第一反应是“先找数据集然后跑个模型”。这个顺序其实是反的。金融信贷风控系统的核心不是机器学习算法本身而是先理解业务上“风控”到底在干什么。信贷业务的基本流程是客户提交申请平台采集客户信息收入、负债、征信记录、历史借贷行为然后决定是否放款、放多少、利率定多少。风控要做的事情就是在这个流程里尽可能准确地回答一个问题这个客户未来会不会违约。所谓“违约”在业务上需要明确界定。通常的做法是定义表现期比如观察客户从放款日开始90天内是否出现逾期超过30天或90天的情况。逾期30天以内的往往是短期遗忘或资金周转问题逾期90天以上基本可以认定为坏客户。这个好坏样本的界定直接影响后续建模因为模型学的是“什么样的人会成为坏客户”前提是你先告诉它“谁算是坏客户”。确定了目标之后风控系统要处理的数据就很清晰了客户提交的申请信息、征信报告里的信贷记录、账户行为数据、历史还款流水。这些数据在真实场景里分散在多个数据源格式不一致、字段口径不一致而且量级非常大。一个中等规模的消金平台注册用户几百万到上千万每人的交易流水少则几十条、多则上千条全量数据规模轻松上亿甚至几十亿行。这就是整套系统为什么必须上大数据技术栈的根本原因。1.2 为什么单机MySQL扛不住非得HadoopSpark我见过不少同学在开题报告里写“系统采用Spring Boot MySQL Redis实现”然后被评委一句话问住“数据量多大”如果答案是几千条测试数据那Hadoop和Spark确实没必要出现用Excel就够了。但信贷风控的真实场景不是这样。假设平台有500万存量客户平均每个客户每月产生20条账户流水保留24个月这就是2.4亿条流水记录。再加上征信查询记录、还款计划表、外部数据文件总记录数保守估计在3到5亿条。MySQL单表在千万级时通过索引还能勉强支撑到了亿级复杂聚合查询基本跑不动。更关键的是风控特征计算不是简单按主键查一行数据而是要对每个客户做横向聚合近30天消费总额、近6个月最大逾期天数、历史贷款申请次数、循环贷额度使用率。这些指标需要对全量行为数据按用户维度分组再计算单机数据库在这个量级上很难在合理时间内完成。Hadoop解决的是存储和可靠性的问题。HDFS把大文件切块分布式存储默认3副本任何一台机器挂了数据都不丢Spark解决的是计算效率的问题。它基于内存做计算同一个数据集可以被反复使用而不用每次读磁盘特别适合特征工程这种需要多轮迭代、多次变换的场景。Hadoop和Spark不是替代关系而是分工关系离线数据落HDFSSpark从HDFS读数据做计算计算结果写回HDFS或者进MySQL供上层应用查询。1.3 从源码目录反推整个系统的模块划分拿到一套完整的源码先别急着跑先看目录结构。一套合格的“基于HadoopSpark的大数据信贷风控系统”源码目录结构通常能直接反映系统架构。我按最常见的工程组织方式拆一下模块目录名职责核心技术数据接入ingestion/将外部数据文件上传到HDFS登记元数据HDFS Shell、Java/Scala客户端数据清洗etl/去重、缺失值处理、格式归一化、生成宽表Spark SQL特征工程features/按用户维度计算风控特征输出特征宽表Spark SQL、DataFrame API模型训练model/样本准备、模型训练、评估、保存Spark MLlib决策服务decision/加载模型输出评分执行规则引擎Spring Boot、Python/Java数据可视化web/管理后台展示审批记录、风险分布Vue、ECharts、MySQL这套分层逻辑和真实企业里的大数据风控平台基本一致。源码里如果把“数据处理”和“业务展示”写成一坨答辩时被问到模块边界就很容易露怯。做毕设的同学哪怕代码规模不大也建议保持这个分层因为它既能体现对大数据架构的理解也方便后续单独替换或升级某个模块。2. Hadoop与Spark在风控系统里的分工比你想象得更明确2.1 HDFS海量信贷数据的“底仓”HDFS在整个系统里承担的角色可以理解成一个巨大的、不怕机器损坏的文件仓库。它把文件切分成128MB的块分散存储到集群的不同节点上默认保存3个副本任何一个副本所在节点挂了系统还能从另外两个副本读数据。在一个信贷风控系统里HDFS上存放的数据一般是这样的数据文件内容建议数据量演示用application.csv客户申请记录含收入、负债、工作年限等200万条billing.csv账单流水含消费金额、消费时间等6000万条repayment.csv历史还款记录含还款日期、还款状态3000万条credit_query.csv征信查询记录含查询机构、查询原因1000万条你需要在HDFS上先建立一套目录来管这些数据。常用的命令是hdfs dfs -mkdir -p /user/hadoop/credit/raw hdfs dfs -mkdir -p /user/hadoop/credit/cleansed hdfs dfs -mkdir -p /user/hadoop/credit/feature hdfs dfs -mkdir -p /user/hadoop/credit/model hdfs dfs -put application.csv /user/hadoop/credit/raw/这个目录划分不是为了好看而是让数据流转路径清晰可追溯原始数据进raw清洗后进cleansed特征表进feature训练好的模型文件进model。答辩的时候评委如果问“你这套系统的数据从哪来到哪去”直接把目录结构说清楚就是最好的回答。2.2 YARN计算资源怎么分给“干活的人”存储层搞定之后接着要回答的问题是Spark作业跑在集群上资源到底怎么分配。YARN就是Hadoop生态里的资源调度器它统一管理集群里每台机器的CPU和内存谁提交了作业YARN就给它分一个ApplicationMaster然后由ApplicationMaster向ResourceManager申请Container来跑Executor。Spark on YARN有两种部署模式client模式和cluster模式。毕设阶段我用得最多的是cluster模式因为Driver也运行在集群里提交作业的机器退出后任务不受影响。一个典型的提交命令长这样spark-submit \ --master yarn \ --deploy-mode cluster \ --class com.example.credit.ETLJob \ credit-system-1.0.jar \ /user/hadoop/credit/raw/application.csv \ /user/hadoop/credit/cleansed/application_clean.parquet理解YARN的意义在于评委经常顺着“Spark作业怎么跑起来”往下问“那集群资源怎么分配的”。你不必背很多参数能说清楚“Spark的Driver和Executor都是运行在YARN的Container里由ResourceManager统一调度”就已经赢了大多数只顾着调代码的同学。2.3 Spark SQL特征计算的生产线信贷风控的特征计算是整个系统里代码量最大、最能体现“大数据处理能力”的部分。大部分风控特征本质上就是“对用户历史行为做聚合统计”。比如“近30天消费总金额”就是对billing表按用户ID分组、筛选近30天数据、求和“历史最大逾期天数”就是对还款记录按用户ID分组、取最大值。这些操作如果用MapReduce写过程非常痛苦——每个聚合逻辑都要写Mapper和Reducer而且中间结果反复落磁盘。用Spark SQL同样的逻辑只需要读表做SQL变换。我拿“统计用户近30天消费特征”举个例子from pyspark.sql import SparkSession from pyspark.sql.functions import sum, max, count, col spark SparkSession.builder \ .appName(feature_engineering) \ .config(spark.sql.shuffle.partitions, 200) \ .getOrCreate() df spark.read.parquet(/user/hadoop/credit/cleansed/billing_clean.parquet) feature_df df.filter(col(bill_date) 2024-01-01) \ .groupBy(user_id) \ .agg( sum(bill_amount).alias(total_amount_30d), count(bill_id).alias(total_count_30d), max(bill_amount).alias(max_amount_30d) ) feature_df.write.mode(overwrite) \ .parquet(/user/hadoop/credit/feature/user_consume_feature.parquet)这段代码读起来很简单但它背后的逻辑是分布式的数据被按用户ID分不到不同节点每个节点并行计算自己那部分分组最后再汇总。Spark SQL的Catalyst优化器会自动做谓词下推和列裁剪也就是说它只读业务需要的列和行并不是把全表扫描一遍以后才做过滤。这个特点在答辩时可以重点讲因为这直接回应了“为什么比直接写Python跑循环快”的问题。2.4 MLlib模型训练与评估的统一入口特征算完之后模型训练就顺理成章了。Spark MLlib提供了一套完整的机器学习Pipeline接口用起来和sklearn思路相近但底层能处理更大规模的数据。信贷风控最常用的模型是逻辑回归。不是因为逻辑回归预测最准而是因为它在信贷这种强监管场景下可解释性最好。贷款机构被客户投诉“为什么拒绝我的申请”时必须能给监管方解释清楚原因逻辑回归的系数天然就是这个解释。一段基础的MLlib训练流程如下from pyspark.ml.feature import VectorAssembler, StandardScaler from pyspark.ml.classification import LogisticRegression from pyspark.ml.evaluation import BinaryClassificationEvaluator data spark.read.parquet(/user/hadoop/credit/feature/feature_table.parquet) # 注意label列是1表示坏客户0表示好客户 feature_cols [x1, x2, x3, x4, x5, x6, x7, x8] assembler VectorAssembler(inputColsfeature_cols, outputColfeatures_raw) scaler StandardScaler(inputColfeatures_raw, outputColfeatures) lr LogisticRegression(featuresColfeatures, labelCollabel) from pyspark.ml import Pipeline pipeline Pipeline(stages[assembler, scaler, lr]) train_df, test_df data.randomSplit([0.8, 0.2], seed42) model pipeline.fit(train_df) pred_df model.transform(test_df) evaluator BinaryClassificationEvaluator(labelCollabel, metricNameareaUnderROC) auc evaluator.evaluate(pred_df) print(fAUC: {auc})训练完的模型可以保存到HDFS指定目录比如model.write().overwrite().save(/user/hadoop/credit/model/lr_model)。决策服务启动时再把这个模型加载回来对新的客户特征做预测打分。训练和预测分离这个设计在毕设源码里如果体现出来了答辩评价会明显上一个档次。3. 信贷风控链路在代码层面是怎么流转的从原始数据到审批分数3.1 第一步永远是洗数据脏数据不处理后面全白搭任何真实场景下的数据都不会是干净的。信贷数据尤其严重我见过的情况包括收入字段填了0或者负数、年龄字段出现200岁、同一个用户在同一天产生两条完全一样的流水、征信查询日期格式五花八门。这些脏数据不处理干净后面做特征计算的时候会得出荒谬的结果比如把月收入为0的用户和月收入50万的用户混在一起训练模型很容易被异常值带偏。数据清洗阶段的Spark作业一般做四件事去重、缺失值处理、异常值过滤、格式标准化。有一个原则是“保留痕迹、可追溯”就是清洗前后数据量要做对比知道每条规则过滤了多少数据。这不只是工程习惯也是论文里“数据分析”章节的重要素材。一个典型的清洗逻辑from pyspark.sql.functions import col, when, to_date df spark.read.csv(/user/hadoop/credit/raw/application.csv, headerTrue, inferSchemaTrue) # 1. 去除完全重复的记录 df df.dropDuplicates([user_id, apply_date]) # 2. 收入字段负数视为缺失缺失值填充为该收入分组的众数或中位数 df df.withColumn(monthly_income, when(col(monthly_income) 0, None).otherwise(col(monthly_income))) # 3. 过滤明显异常年龄不在18到70之间的记录直接剔除信贷产品一般服务这个年龄段 df df.filter((col(age) 18) (col(age) 70)) # 4. 日期统一为yyyy-MM-dd格式 df df.withColumn(apply_date, to_date(col(apply_date), yyyy-MM-dd)) df.write.mode(overwrite).parquet(/user/hadoop/credit/cleansed/application_clean.parquet)这个过程看起来基础但决定了下游特征计算的上限。我在指导毕设时经常说一句话模型的结果上限由数据质量决定特征工程决定模型能不能接近这个上限算法反而是最后才需要考虑的事情。3.2 特征工程把“还款能力”翻译成机器能算的数字模型不认“收入挺高的”“征信不太好”这种模糊描述它只认数字。特征工程就是把业务经验翻译成数值特征的过程。信贷风控里常用的特征可以分成几类。基本信息类年龄、工作年限、月收入、学历等负债类负债率月还款总额除以月收入、信用卡已用额度、循环贷额度使用率等行为类近30天消费笔数、近30天最大单笔消费、近6个月征信查询次数等历史表现类历史最大逾期天数、历史贷款笔数、历史结清率等。每个特征怎么定义需要和业务一一对应。以“负债率”为例它的含义是客户每月的债务支出占收入的比例。如果月收入1万、月还款8000负债率就是80%风险明显偏高。这个特征在Spark里的计算逻辑是repay_monthly spark.read.parquet(/user/hadoop/credit/cleansed/repayment_clean.parquet) \ .groupBy(user_id) \ .agg(sum(monthly_payment).alias(total_monthly_payment)) income spark.read.parquet(/user/hadoop/credit/cleansed/application_clean.parquet) \ .select(user_id, monthly_income) feature_x1 income.join(repay_monthly, user_id, left) \ .withColumn(x1_debt_ratio, col(total_monthly_payment) / col(monthly_income))特征列命名建议直接用x1、x2这种编号配套一个特征字典表说明每个编号的含义。为什么这么做因为很多信贷比赛和真实项目都是用这种脱敏字段命名既保护了业务敏感信息又方便代码统一处理。答辩时评委问“x3是什么”你能立刻对答出来说明你对特征体系是清楚的。3.3 样本不均衡与模型评估别让准确率骗了你训练风控模型时几乎一定会遇到一个问题坏样本太少。真实信贷场景中违约客户比例通常低于5%有些产品甚至不到1%。如果直接把原始数据丢给逻辑回归训练模型会学到“所有客户都是好的”整体准确率哪怕做到97%也没有任何实际意义因为那97%全是被正确分类的好客户而真正需要识别的坏客户一个都没抓住。处理样本不均衡常见做法有几种。下采样从好客户里随机抽让好坏比例达到1:1或者2:1训练速度快但会浪费大量好客户数据上采样用SMOTE等方法人工合成少数类样本效果通常更好但实现稍复杂还有一种思路是调整逻辑回归的classWeight参数给坏客户更高权重这样模型会被“逼着”更关注少数类。评估指标也不能只看准确率。信贷风控最常用的两个指标是AUC和KS。AUC衡量的是模型把好客户排在坏客户前面的能力0.5等于瞎猜0.7以上算有区分度0.8以上已经是很不错的模型。KS的直观含义是模型能区分出好坏样本的最大差异一般超过0.3就可以投入业务使用。举个例子如果测试集有1万条样本模型AUC是0.78、KS是0.42这个结果在毕设里已经是可以拿得出手的水平。关键是要在答辩PPT里把混淆矩阵、坏客户召回率这些指标一起放出来证明你理解“准确率骗人”这件事。3.4 决策引擎分数出来后谁来拍板模型输出的原始结果是“违约概率P”但实际业务系统一般不会直接用P来做判断而是把它映射成一个标准化的分数。业界最经典的映射方式是评分卡Score offset factor * log(odds)其中odds是好客户概率除以坏客户概率offset和factor是两个常数控制分数的基准和刻度。比如设定“当坏好比为1:1时分数为600每增加20分坏好比翻倍”那么公式中的factor和offset可以通过方程组计算出来。评分算出之后决策环节通常是几张规则的组合。用伪代码表示就是def make_decision(score, blacklist_flag, debt_ratio): if blacklist_flag: return REJECT if score 700: return APPROVE if score 600: return REJECT if debt_ratio 0.5: return MANUAL_REVIEW return MANUAL_REVIEW注意决策是“规则模型”的组合不是单一模型分数一票定生死。黑名单命中直接拒绝、高负债即使分数不低也要人工审核这些都是业务层面的硬性约束模型解决不了。源码里把“决策逻辑”独立成一个模块而不是写死在模型训练代码里是这套系统专业性的重要体现。4. 答辩时评委老师最常追问的几个技术点提前想清楚4.1 为什么离线分析用Spark而不是纯MapReduce这个问题几乎是必问的。MapReduce和Spark都能做分布式计算但MapReduce的每个计算步骤都要把中间结果写到磁盘下一个步骤再从磁盘读回来。一次简单的分组聚合可能要经历多轮写盘和读盘在数据量大、迭代次数多的情况下大量时间都花在磁盘I/O上。Spark不一样它基于内存计算RDD在内存中以分区形式保存同一个RDD可以被多个操作重复使用中间结果不用频繁落盘。尤其是在特征工程这种“读一次数据、然后反复做变换”的场景Spark的性能优势非常明显。如果只做一次性的简单统计比如对全量数据算个平均值MapReduce也够用但风控特征计算要对同一份用户行为数据做几十个维度的聚合Spark明显更适合这种迭代式计算需求。4.2 “你的数据量多大”——大数据系统最怕这个问题评委问这个问题主要想看你的Hadoop和Spark是不是真的有必要。如果数据量只有几千条、文件总大小只有几MB那整个架构就是“为大数据而大数据”。建议在毕设阶段把模拟数据量做到一个能撑起架构的量级。比如生成200万条申请记录、6000万条账单流水、1000万条查询记录压缩后HDFS上数据总量大概在几十GB级别。这样Spark作业跑起来能看到真实的stage耗时、shuffle数据量答辩时能直接展示YARN日志里的执行时间说服力非常强。如果机器配置有限生成不了这么多数据也要在答辩时给出理论估算“系统设计支撑千万级用户、亿级流水本实验环境用300万条数据完成全链路验证。”这比直接说“数据量不大”体面得多。4.3 模型可解释性怎么保证风控模型不能只看准确率信贷风控是强监管领域银行和消金公司使用模型必须符合可解释性要求。客户的贷款申请被拒绝后平台需要能给出合理解释监管机构会查。如果用一个纯黑盒模型模型内部逻辑完全不可解释业务上是很难落地的。这也是为什么逻辑回归在信贷场景里占据重要地位。每个特征都对应一个权重系数系数为正说明该特征值越大坏客户概率越高反之亦然。比如“历史逾期次数”这个特征系数是0.8那说明逾期次数越多客户违约概率越高这个解释既直观又符合业务常识。如果用了随机森林或XGBoost这类集成模型也要通过特征重要度排序或者SHAP值来解释每个特征对预测结果的贡献。答辩时可以现场展示特征重要性排序图让评委看到你在意的不只是模型效果还有业务层面的可用性。4.4 特征工程里最容易被问倒的坑特征工程相关的提问最阴险的一个角度是“数据泄漏”。所谓数据泄漏就是训练模型时用了未来才能知道的信息导致训练时模型表现很好、上线后效果崩塌。举个典型例子建模目标是要预测客户未来90天内会不会逾期但特征里包含了“客户在第60天是否已经发出还款提醒”。这个特征和预测目标高度相关训练时AUC可能高达0.95可它本质上是事后才知道的信息上线时会失效甚至违法操作。正确做法是严格切分观察期和表现期观察期提取特征表现期定义好坏标签两者时间上不能有交集。还有一个常见坑是随机切分数据集导致同一客户既出现在训练集又出现在测试集。银行客户可能有多次贷款记录同一个客户的多条样本如果被切到两边模型相当于考试时见过答案。正确做法是按客户维度切分保证训练集和测试集的客户完全不重合。4.5 跨数据源的JOIN和倾斜问题怎么处理特征工程里大量涉及JOIN操作客户表join账单表、账单表join还款表。当一个超大表和一个比较小的维表做JOIN时Spark默认会走Shuffle Join也就是把两边数据都按关联键重分布到各个节点这个过程会产生大量网络传输和磁盘写操作。优化方法是用广播变量。当小表足够小默认10MB以下可以调大时把维表广播到每个Executor节点的内存里大表在本地直接关联根本不用shuffle。在Spark代码里可以用broadcast函数显式指定from pyspark.sql.functions import broadcast result large_df.join(broadcast(small_df), user_id, left)另外要注意数据倾斜。所谓数据倾斜就是某个key对应的数据量特别大导致所有数据都压到同一台机器上计算其他节点空闲等着。信贷数据里很常见的就是几个头部用户有几十万条流水而大多数用户只有几条。缓解思路包括加盐给热点key加随机前缀再拆散、提高shuffle分区数、或者过滤掉极端热点的用户单独处理。能把这几个词说清楚评委基本就能断定你的Spark确实是实操过的不是看几篇博客就能编出来的。4.6 有没有考虑近实时化从离线到实时的演进如果评委问“现在是离线批次计算能不能做到实时代审批”这并不代表你的毕设不完整而是想看你有没有架构演进意识。离线和实时是两种不同场景。当前系统用Spark批处理计算特征适合每天定时跑批、给第二天的审批策略提供依据。要做到秒级实时审批需要引入Kafka接收业务系统产生的实时事件用Spark Structured Streaming或Flink做流式计算实时更新客户特征再用Redis缓存高频访问的决策规则模型服务化后提供在线预测API。毕设阶段不用真的实现完整实时链路但能在答辩里清晰描述“如果要做实时审批我会怎么设计”这比死记硬背概念要加分很多。毕竟系统的演进方向本身就是评委考察你知识广度的方式。5. 从“能跑通”到“能得高分”实操中容易踩的几个坑5.1 环境搭建伪分布式到三节点别在这上面耗太久第一个坑往往不是代码而是环境。很多同学第一步就卡在Hadoop集群搭建上反复折腾SSH免密、Java环境变量、端口冲突两三个星期过去了项目还没开始写。我的建议是先把伪分布式跑通再考虑多节点。伪分布式模式下HDFS、YARN、Spark都装在一台机器上虽然不体现分布式能力但足够用来调试代码和验证流程。实验环境内存小于8GB的话老老实实用伪分布式内存16GB以上可以用三台虚拟机搭一个简单集群每台机器5GB内存留给Hadoop和Spark剩下的留给操作系统。几个容易踩的细节点JDK版本必须和Hadoop、Spark版本匹配Hadoop 3.x需要JDK 8以上但不要用太高版本ssh免密登录要确保能从客户端免密登录到所有DataNodehdfs-site.xml里副本数如果只有一个节点就配置为1否则数据会一直处于“复制中”状态yarn-site.xml里要给NodeManager的可用内存设一个明确值不然YARN可能把所有内存都吃掉导致系统卡死。5.2 HDFS小文件问题碎文件太多NameNode要炸HDFS擅长存大文件小块数据写入会产生大量元数据NameNode内存会被撑爆Spark计算时也会因为FileInputFormat按文件生成切片而产生过多task每个task又要启动和调度效率极低。前面把所有清洗结果都写成Parquet文件并且用mode(overwrite)如果不控制分区数默认每批次可能产出几十个小文件。解决办法很简单每次写出前用repartition(1)或coalesce(1)合并分区。对于演示项目每个中间结果尽量控制在少量大文件比如一张特征宽表写成一个Parquet文件后续读取更快也方便查看。df.repartition(1).write.mode(overwrite).parquet(/user/hadoop/credit/feature/feature_table.parquet)5.3 Spark资源参数不是随便填的Spark作业报错最多的原因不是代码逻辑而是内存不够。spark-submit里如果不对资源参数做规划默认配置经常会导致Executors频繁GC甚至OOM尤其是特征工程阶段做大量JOIN和聚合的时候。一个适合3节点、每台16GB内存实验集群的参数模板大致如下spark-submit \ --master yarn \ --deploy-mode cluster \ --driver-memory 4g \ --executor-memory 8g \ --executor-cores 2 \ --num-executors 3 \ --conf spark.sql.shuffle.partitions200 \ --class com.example.credit.FeatureJob \ credit-system-1.0.jar这里的逻辑是每台机器物理内存是16GB要给操作系统和NodeManager预留3到4GBYARN能分配给Container的内存大概在11GB左右所以一个Executor给8GB比较安全剩下的留给ApplicationMaster和系统开销。spark.sql.shuffle.partitions控制shuffle后的分区数默认值是200数据量大时可以调到400或者800避免单分区数据量过大导致内存溢出。调整参数后你可能发现作业总耗时反而下降了这种“调优反而更快”的经验写在论文里是非常有含金量的。5.4 仿真数据生成最容易翻车的环节很多同学一开始就想着“找真实数据”但真实信贷数据涉及隐私网上很难拿到。于是有人直接用Python的random函数随便造几百万条数据结果模型AUC只有0.55和瞎猜差不多然后就开始怀疑算法、调参折腾半天没效果。问题往往出在数据上完全随机生成的特征之间没有任何相关性逻辑回归当然学不出规律。正确的做法是按业务规律造数据。好客户和坏客户要服从不同分布比如好客户月收入均值在15000左右坏客户月收入均值在5000左右好客户负债率大多在0.3以下坏客户负债率经常超过0.6好客户历史逾期次数集中在0到1次坏客户集中在3次以上。让不同样本在特征上天然存在区分度模型训练时才能学到有效规律。如果想让数据更像真实的还可以设置特征间的相关性收入越高信用卡使用率越低年龄越大工作年限越长。不要小看造数和数据分布的合理性它能直接决定你模型评估结果的好看程度也是论文里“实验数据说明”部分能不能让评委信服的关键。最后再说点实际的这套系统做完之后我最大的体会是大数据风控毕设的真正分水岭不在算法选型而在数据工程。你可以把逻辑回归换成随机森林再把随机森林换成XGBoost但如果HDFS上的数据是脏的、特征口径是错的、样本是失衡的、资源参数是瞎配的任何模型都救不回来。评委老师见多识广一问数据量、二问特征怎么算、三问模型怎么评估、四问资源怎么调你有多少水分当场就能试出来。来源码里最值得反复看的是三个地方原始数据清洗作业写得规不规范特征口径和业务含义对不对得上决策模块有没有把规则和模型分开。这三个地方能讲清楚你答辩的心态都会不一样。还有一个小技巧提交Spark作业的时候打开--verbose日志按stage记录每个阶段的耗时答辩展示时把YARN日志里的执行时间截图放在PPT里比说一百句“这个系统能支撑海量数据”都有说服力。数据工程这件事很难写出花来但恰恰是它决定了你的项目是“作业”还是“作品”。本文还有配套的精品资源点击获取
返回列表