Spark ML全流程实战:从数据处理到模型部署的完整指南

发布时间:2026/9/9 4:55:26
Spark ML全流程实战:从数据处理到模型部署的完整指南 年初接了一个信贷场景的活要基于几千万行历史数据跑违约预测模型。一开始图省事打算照老路子用Pandas加sklearn先跑通结果连数据都读不进去内存直接被打满。后来把整条链路切到Spark上从数据预处理、特征工程、模型训练到部署一口气串下来才算把这个项目啃完。这篇文章就是把当时那套完整流程做一个复盘。内容围绕Spark ML的DataFrame API展开从环境准备讲起再讲数据预处理、Pipeline建模、调参评估最后落到生产环境里模型怎么部署、上线后会踩哪些坑。整个流程如果你只用本地小数据集练习过机器学习会很有参考价值如果你正好在Spark集群上跑过ETL但没碰过ML也适合用来把两边的知识串起来。1. 为什么用Spark做机器学习是一条绕不开的路1.1 传统单机方案的瓶颈在哪先说一个很现实的问题为什么不用Pandas加sklearnPandas处理百万行以内的表格确实顺手但一旦数据量到千万、亿级它的内存占用会非常夸张。一个3000万行、300列的数据集就算字段全是数值型在Pandas里动辄占用20到30GB内存读CSV的时候还可能翻倍。在实际项目里数据往往还带字符串、时间、文本字段内存压力更大。很多人第一反应是“加大内存”。但内存加到128GB之后会发现另一个问题单机训练时间不可控。一次交叉验证要在几百棵树里做特征分裂单机要跑到天亮。每一次调参都等于重跑一轮这种迭代节奏根本没法支撑项目交付。另外还有一个隐性成本数据仓库里的数据本身就不在本地。几千万行明细存在Hive或者云数仓里你用Pandas还得先导出来光导出和转换就够折腾。Spark的定位是“面向大规模数据处理的统一引擎”它能直接在分布式存储上做计算数据和计算不分离这在实际生产中比单机方案省掉很多搬运工作。1.2 Spark ML的定位与新旧API选型Spark里的机器学习库经历了两个阶段。早期叫MLlib基于RDD API写起来更像传统的MapReduce程序后来官方推出基于DataFrame的spark.ml包也就是现在讲的Spark ML。新老API最大的区别在于抽象模型不同。RDD API更底层每次操作都在操作分布式集合缺少schema信息很多优化做不了。DataFrame API引入了schema、列式存储和Catalyst优化器Spark能对整条执行计划做优化包括谓词下推、列裁剪、分区裁剪这些在做特征工程时效率差异巨大。所以我的原则很简单新项目一律用spark.ml包不要碰旧的MLlib。如果你在网上搜到一些老代码还在用LabeledPoint和RDD可以关掉不看了那套接口已经不适合生产场景。1.3 Spark ML Pipeline把“流程”变成可复用资产传统机器学习里数据预处理是脚本特征工程是脚本训练脚本最后再单独保存一个模型文件。脚本之间没有统一的结构步骤一多顺序乱了、参数改了、数据对不上了乱成一锅粥。Spark ML的做法是提供一个Pipeline机制。你把数据清洗、特征转换、模型训练这些步骤串成一条流水线流水线本身可以保存、加载、复用。训练时调用fit预测时调用transform整个结构化流程一次成型。这个特点在我们后面做模型部署时价值很大。2. 环境准备本地跑通Spark ML的最小配置2.1 Spark安装和启动的几个关键点先准备环境。你在自己电脑上装个单机版Spark跑通训练流程完全够用。官方目前推荐Spark 3.x配合Java 8/11/17都可以我自己用Java 17时间最长没遇到兼容问题。安装步骤不复杂核心就三步装Java、下载Spark、配置环境变量。下载时直接选Pre-built版本对应Hadoop版本选“without Hadoop”或者带内置Hadoop的包都行本地跑用不到HDFS数据直接读本地文件。export JAVA_HOME/usr/lib/jvm/java-17-openjdk-amd64 export SPARK_HOME/opt/spark-3.5.1-bin-hadoop3 export PATH$SPARK_HOME/bin:$PATH # 本地启动一个交互式环境 pyspark --master local[4]实际提交训练任务别用pyspark交互式环境而是把代码写成.py脚本用spark-submit跑。比如spark-submit \ --master local[4] \ --executor-memory 8g \ train_loan_model.py第一次跑通之后你会发现用Spark跑任务的逻辑其实不复杂真正麻烦的是后面分布式部署和调优。如果需要在集群上跑常见做法是搭一个Standalone集群主节点执行start-master.sh工作节点执行start-worker.sh然后在提交任务时把--master指向主节点的地址。2.2 数据源不只是CSVSpark怎么接Redis中的数据实际业务中需要建模的样本特征常常不全部在数据仓库里。有一回我们需要把Redis中的实时特征和Hive中的历史标签拼接当时第一反应是用Redis客户端逐条查询发现速度完全不可接受。Spark读取Redis一般不是直接有一个内置的spark.read.redis这样一句话而是要引入第三方的Redis连接包或者自己写RDD连接器。但如果你只想做一次性建模有一个成本更低的方案用Spark直接读取一份从Redis批量导出的数据文件。# 方案A读取批量导出后的数据 redis_export spark.read.json(hdfs:///tmp/redis_feature/) # 方案B使用自定义分区逐个读取Redis def read_redis_partition(partition): import redis r redis.Redis(hostredis-host, port6379, decode_responsesTrue) for key in partition: yield (key, r.get(key))方案A在生产中我会优先推荐。原因很简单建模之前要做的是离线批量特征拼接不是在线事务查询把所有Redis键一次性批量导出比逐条调用效率高几个数量级。等模型上线后需要实时特征再做服务化查询那是部署阶段的事了。2.3 加载数据并建立基线认知数据加载使用DataFrame API格式支持很全。我们以经典的信贷违约预测数据为例这类数据集一般包含借款人历史借贷记录、还款状态、金额等信息列名类似RevolvingUtilizationOfUnsecuredLines、age、NumberOfTime30-59DaysPastDueNotWorse标签列是SeriousDlqin2yrs。from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(LoanDefaultPrediction) \ .getOrCreate() df spark.read.csv( data/training.csv, headerTrue, inferSchemaTrue ) df.printSchema() df.describe().show()拿到数据第一件事不是建模而是花时间搞清楚每一列的业务含义确认数据形态跟你的预期是否一致。不要急着写代码做预处理先把字段列出来检查空值率、最大值、最小值建立对整个数据质量的直观感受。这一步做得越细后面建模时返工越少。3. 数据预处理最花时间的那80%3.1 缺失值处理要结合业务不只是填个均值信贷数据里缺失值非常常见比如用户没有某个信贷产品对应的额度字段就为空。处理缺失值有几种主流思路实际用哪个要结合字段含义删除缺失比例过高的列比如超过60%就直接剔除用均值、中位数或众数填充适合分布比较正常的数值型字段用前值或后值填充适合时间序列相关的特征把缺失单独变成一个类别适合业务上“缺失本身有意义”的情况。Spark MLlib里提供了Imputer支持用均值或中位数填充。一个容易忽略的细节是填充统计量必须从训练集计算然后用同样的数值去填充测试集。这个Spark已经帮你考虑到了——Imputer.fit会在训练集上算出填充值模型保存后再transform时用的就是训练时的统计量。from pyspark.ml.feature import Imputer imputer Imputer( inputCols[MonthlyIncome, NumberOfDependents], outputCols[MonthlyIncome_imp, NumberOfDependents_imp] ).setStrategy(median) df_imp imputer.fit(df).transform(df)建议不要把MonthlyIncome_imp直接替换掉原始列保留原列作为“缺失指示特征”可能有额外信息。先新增一列MonthlyIncome_missing标记是否为空再交模型去拟合这个交互关系往往效果更好。3.2 分类型变量编码的隐藏风险分类型变量要转成数值最直接的方式是StringIndexer加OneHotEncoder的组合。StringIndexer的作用是把字符串类别映射成数字索引。有个容易踩坑的点是它默认按类别出现频率排序而不是按字母排序。这意味着如果你在训练集上fit测试集出现了训练集中没有的新类别程序会直接报错。解决方案是在StringIndexer里设置setHandleInvalid(keep)给未知类别单独保留一个映射。字符串转索引之后还要再做一步OneHotEncoder。原因是直接给模型喂整数编号会引入顺序关系类别3和类别1的数值距离会被模型理解为2。但分类型变量的不同取值之间没有顺序含义OneHot可以把每个类别映射成独立维度。如果类别本身确实有顺序比如“低中高”那直接保留索引列可能更合理。from pyspark.ml.feature import StringIndexer, OneHotEncoder indexer StringIndexer( inputColcategory, outputColcategory_idx, handleInvalidkeep ) encoder OneHotEncoder( inputCols[category_idx], outputCols[category_vec] )OneHotEncoder输出的是一个稀疏向量。数据量大时这种存储方式非常省内存大多数维度是0只记录非0位置。3.3 特征组装和标准化注意别把顺序搞反数值型和分类型特征处理完最终要把所有特征合并成一个features向量列这一步使用VectorAssemblerfrom pyspark.ml.feature import VectorAssembler, StandardScaler asm VectorAssembler( inputCols[ MonthlyIncome_imp, MonthlyIncome_missing, DebtRatio, age, NumberOfDependents_imp, category_vec ], outputColraw_features ) scaler StandardScaler( inputColraw_features, outputColfeatures, withStdTrue, withMeanTrue )顺序很重要一定要先向量化再标准化。不能用标准化的方式去挨个处理原始列那样在分布式系统里既要额外维护每列的统计量又容易在组装时出错。需要提醒的是withMeanTrue会把数据做中心化在稀疏向量上会导致向量变稠密极度消耗内存。如果特征维度很高且大部分为0建议withMeanFalse只做方差缩放或者使用SparseStandardScaler等更合适的策略。3.4 划分数据集时最容易犯的错很多人在整个DataFrame上直接做标准化、填充、编码然后才切分训练集和测试集。这个做法存在数据泄漏的风险。你填充平均值时已经用到了测试集的数据信息评估出的指标会比真实表现高上线后模型效果立刻打折。正确的做法是先用randomSplit把数据拆成训练集和测试集然后所有的fit都只在训练集上执行。Pipeline机制能帮你强制做到这一点把处理步骤放进去最后一次性拟合。train_df, test_df df.randomSplit([0.8, 0.2], seed42) print(train_df.count(), test_df.count())4. 用Pipeline把整个训练流程规范起来4.1 Pipeline的设计哲学机器学习里大致有两类组件。一类是Transformer它有一个transform方法作用是把一个DataFrame变成另一个DataFrame比如Imputer、VectorAssembler、训练好的模型。另一类是Estimator它有一个fit方法作用是在DataFrame上训练并产出Transformer比如逻辑回归、随机森林。Pipeline的任务就是把这两类组件按顺序串起来最末尾通常是一个Estimator。调用pipeline.fit(train_df)时Spark会依次执行每一步中间步骤直接调用transform最后一个Estimator调用fit产出模型。整个过程结束后返回一个PipelineModel。这个机制最大的价值是杜绝手工步骤遗漏。有一次我们线上模型和离线回测效果差异巨大排查了一整天才发现预处理少跑了一步标准化这种错误以后再也没有出现过因为整个链条都被固化在同一个Pipeline里了。4.2 一个可以直接抄的Pipeline代码把前面讲的预处理步骤串起来写成一个流水线训练一个随机森林分类器from pyspark.ml import Pipeline from pyspark.ml.classification import RandomForestClassifier from pyspark.ml.evaluation import BinaryClassificationEvaluator # 第一步缺失值填充 imputer Imputer( inputCols[MonthlyIncome, NumberOfDependents], outputCols[MonthlyIncome_imp, NumberOfDependents_imp] ).setStrategy(median) # 第二步字符串索引与独热编码 indexer StringIndexer( inputColcategory, outputColcategory_idx, handleInvalidkeep ) encoder OneHotEncoder( inputCols[category_idx], outputCols[category_vec] ) # 第三步特征向量化 assembler VectorAssembler( inputCols[ MonthlyIncome_imp, DebtRatio, age, NumberOfDependents_imp, category_vec ], outputColraw_features ) # 第四步模型训练 rf RandomForestClassifier( featuresColraw_features, labelCollabel, numTrees100, maxDepth8, seed42 ) # 组装流水线 pipeline Pipeline(stages[imputer, indexer, encoder, assembler, rf]) # 在训练集上执行端到端训练 model pipeline.fit(train_df) # 在测试集上一键完成与训练时完全一致的预处理并预测 predictions model.transform(test_df)注意整个训练过程没有手动去调用每个转换步骤预处理逻辑和模型是一体的。预测的时候只需要一行transformSpark会自动按保存的顺序加载预处理规则并应用到新数据上这在部署阶段极其重要。4.3 参数调优时别让分布式框架变成慢动作训练完一个初始模型要调参。Spark提供的CrossValidator跟sklearn的用法类似但它会把每组参数都跑一遍完整的分布式训练绝对时间比单机要久。一个包含几十组参数组合的交叉验证任务跑几个小时很正常。经验是先小规模探索。第一轮网格不要太大每组参数选择相隔较远的取值比如numTrees取20和100maxDepth取5和10先找出大致方向。第二轮再根据结果围绕较优区域精细搜索。如果数据实在太大可以先在一个采样子集上跑参数实验锁定候选参数后再全量训练。from pyspark.ml.tuning import CrossValidator, ParamGridBuilder grid ParamGridBuilder() \ .addGrid(rf.numTrees, [20, 50, 100]) \ .addGrid(rf.maxDepth, [5, 10]) \ .build() evaluator BinaryClassificationEvaluator( metricNameareaUnderROC ) cv CrossValidator( estimatorpipeline, estimatorParamMapsgrid, evaluatorevaluator, numFolds3, seed42 ) cv_model cv.fit(train_df)CrossValidator输出的cv_model本身也是一个PipelineModel可以直接用来预测。不过要小心它会保留所有训练出来的子模型占用较大存储空间。生产环境确认好参数后用最优参数重新训练一次只保存最终模型就够了。5. 模型评估与分布式训练中的细节5.1 评估指标要跟着业务场景走信贷违约预测常用的指标是AUC和KS。AUC衡量的是模型把正样本排在负样本前面的能力对样本不平衡不敏感适合做模型排序能力评估。KS更直观地体现了在某个分数切点下能识别出多少坏客户在风控场景里常被作为审批策略的参考指标。from pyspark.ml.evaluation import BinaryClassificationEvaluator evaluator_auc BinaryClassificationEvaluator( labelCollabel, metricNameareaUnderROC ) auc evaluator_auc.evaluate(predictions) print(fTest AUC: {auc})还有一个容易忽略的问题类别不平衡。如果坏样本占比只有1%模型完全可以靠预测全部为“好”来获得很高的准确率但这样的模型没有意义。常用策略是对少数类做加权Spark里的分类器大多支持weightCol参数也可以直接设置classWeight让模型更关注少数类样本。5.2 资源参数设不对训练性能天差地别在集群上训练资源参数没配好会非常难受。常见的两个问题一是executor内存太小导致OOM二是分区数太多导致任务调度开销过大。我习惯用这种配置作为起点spark-submit \ --master yarn \ --deploy-mode cluster \ --num-executors 10 \ --executor-cores 4 \ --executor-memory 8g \ --driver-memory 4g \ --conf spark.sql.shuffle.partitions200 \ train_loan_model.pyspark.sql.shuffle.partitions决定shuffle时的分区数。一个经验值大约是每个executor核心数的2到3倍集群总核心数100左右时设成200比较稳妥。设得太小会导致单个任务处理数据量太大处理太慢设得太大则任务调度开销会拖累整体速度。随机森林在Spark里是按树并行训练的每棵树之间没有依赖关系所以增加executor数量对随机森林的训练速度提升非常明显。反过来像逻辑回归这类需要多轮迭代的算法网络通信开销占比高单纯加机器不一定线性加速。5.3 特征处理阶段要留意数据倾斜数据倾斜在做特征拼接时经常出现。按用户ID做join时如果少数头部用户的记录量非常大某些分区要处理的数据远超其它分区整个任务就会卡在少数几个task上。排查方法很简单。Spark UI里能看到每个stage的task耗时分布如果少数task跑了几十分钟其它几十秒就结束大概率倾斜了。解决思路通常是给join key加盐或者先用聚合把数据打散。如果倾斜的key本身在业务上意义不大直接过滤掉也是方案。6. 模型部署到生产的几条路线6.1 离线批量预测直接复用PipelineModel如果业务场景是每天跑批比如每天凌晨对一批新申请用户做风险评分最简单的部署方案就是让Spark任务定时加载训练好的PipelineModel对新数据做transform。from pyspark.ml import PipelineModel # 加载训练时保存的完整流水线 loaded_model PipelineModel.load(model/loan_model) # 应用在新一批数据上 batch_predictions loaded_model.transform(new_orders_df) batch_predictions.select(user_id, probability, prediction).show()这个方案的好处是零改造。预处理逻辑、特征顺序、模型参数全部被打包在同一个目录里不会出现离线训练时和在线预测时特征不一致的问题。在新数据上你只需要保证字段名称与训练时一致即可。要注意的是PipelineModel保存的是一个目录里面包含多份文件需要整体拷贝到目标环境不能只迁移其中某一个模型文件。在分布式文件系统上可以直接保存到HDFS路径如果要在另一套环境加载把目录整个同步过去就行。6.2 在线实时评分PMML/ONNX导出路线如果业务方要求接口实时返回评分每次请求都去起一个SparkContext加载模型显然不现实。这时候会把模型转换成一种脱离Spark运行的格式最常用的是PMML和ONNX。PMML的全称是预测模型标记语言用XML描述整个数据预处理和模型结构。只要目标服务能解析PMML文件就可以不依赖Spark独立完成推理。相当于把训练阶段的完整特征处理流程连同模型一起“导”出来服务端拿到一条原始数据按PMML内部规则逐步算特征、再算分数。在Spark ML中可以使用JPMML-SparkML这类工具把PipelineModel导出为PMML文件from pyspark.ml import PipelineModel from pyspark2pmml import PMMLBuilder model PipelineModel.load(model/loan_model) pmml PMMLBuilder(spark, model, test_df) \ .buildByteArray() with open(model/loan_model.pmml, wb) as f: f.write(pmml)服务端需要提前确认两件事PMML文件里记录的特征名必须和请求参数名一致PMML对部分算法的支持不如Spark原生全面。实测下来逻辑回归、随机森林这类主流模型导出都比较顺利过于复杂的自定义UDF就需要另想办法。ONNX是另一条路线。Spark ML没有官方ONNX导出组件一般做法是你构建的Pipeline是纯特征工程最终模型用XGBoost或者LightGBM训练然后单独把基学习器导出为ONNX。实际使用中PMML在Spark生态里更顺手所以分享也以这个路线为例。6.3 两种部署方式的对比和选择部署选型依据更多看实时性要求我用一张表来总结场景推荐方式优点缺点每天定时批量打分Spark加载PipelineModel跑批无损、稳定、改动小要维护Spark环境实时接口毫秒级返回PMML/ONNX上标准Java服务响应快、独立部署需解决部分算子兼容问题小流量初上线先用Spark跑批撑住阶段需求最快落地无法支撑高并发在线请求有一个折中方案也经常用训练阶段同时保留Spark PipelineModel作为离线标准答案导出PMML作为在线服务。上线初期随机抽样对比两边打分结果一旦发现偏差可以快速定位是导出环节还是服务端实现的问题相当于给你的部署加了一层安全网。7. 生产环境常见问题和排查笔记7.1 Spark OOM的排查清单跑大规模训练时最常遇到的就是executor OOM。很多人一看到报错就盲目加内存加了还挂其实是没找到根因。我把常见原因整理成了清单数据读取后没有做必要过滤或列裁剪全量字段堆进内存repartition之后分区数太少单分区数据量过大withMeanTrue使稀疏向量变稠密内存翻几倍在transform前做了collect()把海量数据拉回driver端History Server里查看GC时间如果CMS GC占比过高说明内存确实不足要考虑增加executor内存或减少单executor并行度。排查工具优先用Spark UI它能看到每个stage的输入数据量和shuffle读写量。如果明明数据集只有几十GB但shuffle读写了几个TB优先查有没有数据膨胀比如笛卡尔积或者倾斜join。7.2 随机种子带来的复现性问题机器学习的可复现性在分布式环境里比单机麻烦得多。随机森林和决策树在训练时涉及样本采样和特征选择如果不设seed每次结果都会有细微浮动。Spark里大多数模型都暴露了seed参数Pipeline构建时最好统一设置rf RandomForestClassifier( featuresColfeatures, labelCollabel, seed42 )同时randomSplit也要设置seed否则切分出的训练集和测试集每次都不一样比对实验结果就失去了基础。我还习惯在代码开头固定统一的随机种子这样整个流程只要数据不变结果就是可复现的。7.3 部署上线后离线在线不一致的检查模型上线后评分跟离线回测差异大八成出在特征拼接环节。离线训练时你用了一批事后才可以拿到的数据比如“未来30天用户是否逾期”服务端实时请求时根本取不到这个值特征就是错的。这种问题排查很花时间最好是模型立项阶段就和业务方确认清楚每个特征的获取时间点。第二类不一致出在类别编码上。PMML导出后如果类别字典跟训练时不同字符串索引就会错位。解决办法是在导出前对类别做一次排序固定映射关系或者单独保存一份标签字典在服务端启动时加载。7.4 给新手的流程建议整套流程走完你会发现在Spark里做机器学习项目跟平时单机写模型很不一样。数据量越大越要提前花时间在数据分布、特征含义和资源规划上。后端和算法同事沟通时能明确说清楚需要多少资源、多少个分区、预估多久跑完项目推进会顺畅得多。建议上手路径先从kaggle或天池找一份几十万行的结构化数据在自己电脑上把预处理的Pipeline搭起来再换到一百万的量去体会下Spark真正的优势。不要一上来就追求“分布式”“集群”单机先把流程跑通后续加机器只是改一个master地址的事。真正值得反复打磨的是数据质量意识、Pipeline规范以及部署时的特征一致性控制。这个项目做下来我最深的体会有两点。第一机器学习工程的复杂度不在模型本身而在于把数据处理流程从“脚本式的临时操作”变成“可复用、可部署的标准化产物”。第二不管离线效果做得多漂亮只要部署阶段没有保住“训练和推理同源”这条底线迟早会在线上栽跟头。把上面这些坑提前规避掉Spark ML全流程落地这件事真没有想象中那么难。

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询