
在很长一段时间里我处理机器学习任务的第一反应都是把数据抽样到本地然后交给Python那一套生态去跑。直到有一次接手了一个用户流失预测的项目样本量从几万条直接膨胀到几亿条单机训练一次模型要十几个小时特征工程更是慢到让人怀疑人生——从那一刻起我意识到在真正的大数据场景下传统的单机机器学习方案已经到天花板了。也是从那个时候开始我系统地把Spark MLlib用到了实际项目中从用户行为预测到推荐召回从特征处理到模型落地上线算是把这条链路完整地趟了一遍。这篇文章不打算写成一本文档手册的搬运而是想以一个实际使用者的身份把MLlib在大数据AI应用开发中的关键环节、选型理由、踩坑经历和调优经验分享出来。无论你是刚开始接触Spark的初学者还是已经在用Spark做ETL、想进一步尝试分布式机器学习的老手这篇文章的核心目标就是帮你搞清楚一件事什么时候该用MLlib怎么用才能少走弯路。1. 为什么大数据机器学习要选MLlib从一次流失预测项目说起1.1 单机库与分布式计算的分水岭先说说那次流失预测项目的经历。最初我用的是单机版的机器学习框架把用户的登录日志、消费记录、客服交互等特征拼成了一张宽表大概有两百多个维度。数据量在百万级别的时候一切都很美好特征工程跑几分钟就完事模型训练也控制在半小时以内。但随着业务积累数据量很快突破了亿级问题开始集中爆发。单机框架的瓶颈并不在于算法本身而在于数据加载和特征计算。每次训练前要把数据从分布式存储里拉到本地这一项就要消耗几十分钟甚至几个小时。特征处理时涉及的各种关联、聚合、分组统计在单机内存里频繁进行ShuffleOOM成了家常便饭。最痛苦的是我为了跑一个简单的逻辑回归竟然要先写一大堆MapReduce或者Spark ETL脚本把数据处理好再导出来训练整个链路极其割裂。MLlib的出现解决了这个核心痛点它让特征工程、模型训练和模型评估能够无缝衔接在同一个分布式计算框架内完成。数据不需要岀来导出训练算法本身是分布式的数据量从亿级到百亿级只要集群资源够理论上都能跑。从那次项目之后我把所有涉及大规模数据的机器学习任务都迁移到了Spark MLlib再也没有回去过。1.2 MLlib的定位不是用来替代TensorFlow的这是很多初学者容易搞混的地方。MLlib不是万能的它的核心优势在于处理传统的机器学习算法比如逻辑回归、决策树、随机森林、梯度提升树、KMeans聚类、协同过滤等等以及与之配套的特征处理、评估流水线。它并不适合跑深度神经网络——虽然MLlib里保有一个叫MultilayerPerceptronClassifier的多层感知机分类器但它的设计初衷和表达能力显然无法和TensorFlow、PyTorch这套生态相提并论。我的理解是MLlib更像是一个分布式的Scikit-learn解决的是“传统机器学习算法在超大规模数据上的工业化落地”问题。如果你的业务需要的是CTR预估、文本分类、推荐召回这类场景MLlib的线性模型、树模型和协同过滤是非常可靠且高效的。但如果你要做图像识别、自然语言理解、生成式AI那么应该用Spark来做数据预处理和特征工程然后把处理好的数据送给深度学习框架去训练两者配合而不是互相替代。1.3 MLlib的API演进RDD到DataFrame的变迁如果你看过一些老的Spark教程会发现里面充满了RDD[LabeledPoint]这样的代码。那是MLlib早期的写法基于RDD API现在已经逐渐被基于DataFrame的API取代。我强烈建议新项目直接使用DataFrame-based API也就是org.apache.spark.ml这个包下的内容。RDD-based API的问题在于不够直观需要手动构造LabeledPoint这样的对象。和Spark SQL、结构化流之间的集成很差。无法利用Catalyst优化器和Tungsten执行引擎的性能优势。DataFrame-based API的好处是可以直接用SQL或者DataFrame算子做特征工程所见即所得。通过Pipeline机制把特征处理、模型训练串联成一个工作流。保存和加载模型更加标准化。从实际体会来说用DataFrame-based API书写代码的流畅度和可维护性高太多。你可能偶尔在旧博客里看到spark.mllib包下的代码直接无视就好——那个已经是历史遗留了。2. 环境搭建与集群部署这步踩的坑最多2.1 版本选型Spark版本和组件兼容性里藏着魔鬼很多新手一开始没在意版本问题随便下载了一个Spark版本就开始干活结果后面踩到坑一头雾水。MLlib的API在Spark 2.x和Spark 3.x之间有不小的差异尤其是在向量列的处理、正则化参数、交叉验证的并行度这些细节上。我个人建议直接用Spark 3.x以上的版本一方面性能有明显提升另一方面API设计更加成熟。同时要注意Scala版本兼容性。Spark 3.x通常有基于Scala 2.12和Scala 2.13编译的两个版本如果你是用Java或者Scala开发务必保证项目里的Scala版本和Spark发行版保持一致否则运行时会遇到各种奇怪的NoSuchMethodError。用PySpark的话会省心一些不过也要注意Python版本的兼容范围比如Spark 3.4开始已经不支持Python 3.7了。还有Java版本的问题。Spark 3.2之后的版本官方推荐Java 8/11/17但是不同的Spark版本对Java 17的支持完备度不一样。我在Spark 3.3上踩过Java 17的坑某些JVM参数在Java 17里已经废弃导致启动Executor失败。如果你的团队不是特别依赖新版Java的特性建议用Java 8或者Java 11稳妥第一。2.2 本地开发和集群运行的最小配置清单对于本地开发调试不需要一开始就搭集群。最简单的方案是直接下载Spark发行版用本地模式跑。本地模式对于验证代码逻辑、Pipeline设计、调试特征工程非常方便。你可以在代码里这样初始化from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(MLlibDemo) \ .master(local[*]) \ .config(spark.sql.shuffle.partitions, 10) \ .getOrCreate()local[*]表示使用本机所有可用的CPU核心来运行适合开发调试。这里有一个小建议本地模式跑的时候把spark.sql.shuffle.partitions设小一点默认200个分区在本地跑反而会拖慢速度因为每个分区都有调度开销。设成10到20个就够了。生产环境的集群部署如果公司还没有现成的Hadoop集群我建议直接用Spark on YARN的方式。原因很简单YARN可以统一管理计算资源和HDFS的配合最成熟运维体系也有现成的方案。2.3 很多人问的“Spark on YARN CPU只能用1个”问题这个热搜词我太有感触了。很多人在YARN上跑Spark作业时发现明明集群有几十个CPU核心但任务并行度就是上不去每个Executor只分配到了一个vCoreSpark UI里显示的任务数少得可怜。这个问题的根子通常不出在Spark本身而是出在资源调度配置上。YARN的调度器需要知道每个container可以分配多少CPU和内存如果你在提交作业时没有明确设置执行的资源需求YARN只会按照最小的默认配置来分配。举个例子在提交Spark作业时我见过大量这样的命令spark-submit --class com.example.Train \ --master yarn \ --deploy-mode cluster \ --executor-memory 8g \ myapp.jar这里只设置了executor-memory但没有设置--executor-cores和--num-executors。这种情况下Spark会使用默认值通常spark.executor.cores默认为1spark.executor.instances也默认为2。于是你申请了一堆8G内存的Executor但每个Executor只能用1个CPU核心并行度自然上不去。正确的姿势是要同时指定三个关键参数spark-submit --class com.example.Train \ --master yarn \ --deploy-mode cluster \ --num-executors 20 \ --executor-cores 4 \ --executor-memory 16g \ --driver-memory 4g \ myapp.jar这样每个Executor就可以使用4个CPU核心做并行计算。核心要点是YARN上CPU资源闲置的90%情况都是因为作业提交时没有显式地申请核心数。这是我在前面提到的项目里踩过的大坑排查过程花了两三天时间才最终定位到。另外还要注意一个容易混淆的参数spark.default.parallelism。这个参数决定了类似reduceByKey等Shuffle操作的默认分区数但不会影响Executor的核心数或并行度上限。如果Executor核心数不够这个参数设得再大也没有意义。你需要先确认spark.executor.cores和spark.executor.instances再看是否需要调整spark.sql.shuffle.partitions或spark.default.parallelism来优化Shuffle阶段的并发度。3. 分布式环境下的特征工程MLlib的隐形重头戏3.1 DataFrame优先从这个设计原则讲透很多人提到MLlib就想到算法那一层但实际上在真实项目中特征工程的工作量占到了80%以上。在分布式环境下做特征工程首要原则就是:优先使用DataFrame和Spark SQL而不是直接用RDD的map/filter来操作。举个例子假设你需要根据用户的登录日志和订单表来拼接训练特征。在DataFrame的世界里你只需要几行SQL就能搞定train_data spark.sql( SELECT a.user_id, a.age, a.gender, COUNT(b.order_id) AS order_cnt, SUM(b.amount) AS total_amount, AVG(b.amount) AS avg_amount, DATEDIFF(CURRENT_DATE, MAX(b.order_time)) AS days_since_last_order FROM user_profile a LEFT JOIN user_orders b ON a.user_id b.user_id GROUP BY a.user_id, a.age, a.gender )这样的代码不仅可读性强而且Spark SQL的Catalyst优化器会自动帮你做谓词下推、列裁剪等一系列优化。如果你用RDD的map和reduce去做不仅代码冗长易错性能往往也差一个数量级。特征工程完成之后需要把特征列合并为一个向量列供MLlib算法使用。这一步用到的是VectorAssembler它是MLlib里最常用的Transformer之一from pyspark.ml.feature import VectorAssembler feature_cols [age, gender, order_cnt, total_amount, avg_amount, days_since_last_order] assembler VectorAssembler(inputColsfeature_cols, outputColfeatures) train_data_vector assembler.transform(train_data)做完这一步train_data_vector里就多了一个features列这个列就是所有算法的标准输入。3.2 常用的特征处理算子使用场景MLlib里的特征处理算子非常丰富我总结了一下日常高频使用的几个类别对照它们的使用场景算子作用典型使用场景VectorAssembler多列合并为向量把所有特征列合并为模型的输入StringIndexer字符串转数值索引把用户ID、类别ID转为数值OneHotEncoder类别特征独热编码处理无序类别特征StandardScaler标准化连续值特征均值方差标准化MinMaxScaler归一化到[0,1]需要限幅的连续值特征Bucketizer连续值分箱年龄、金额等特征分箱TF-IDF文本向量化文本特征抽取Word2Vec词向量序列特征、文本语义特征PCA降维高维稀疏特征压缩其中的StringIndexer有一个需要注意的细节它会把字符串类别按照出现频次从高到低映射为0、1、2...这样的数值索引但对于新出现的类别默认情况下会报错或设为缺失。如果你期待模型能够处理训练时没见过的类别ID需要在StringIndexer里设置setHandleInvalid(keep)否则上线预测时遇到新ID会直接异常。3.3 特征处理中的隐性问题分布漂移与采样偏差特征工程这一节我想特意多讲一个很多人容易忽略的问题在分布式环境下做标准化或归一化时均值和方差应该是从全量训练集上统计的而不是从某个分区或者抽样数据上计算的。很多初学者会在直接用StandardScaler之前先手动用select(mean(col), stddev(col))去算统计量然后自己写UDF做标准化。这种做法在数据量小时没什么问题但在分布式环境里面临两个麻烦一是手动计算统计量的代码不够优雅二是一旦数据量大反复扫描数据的开销很大。MLlib的StandardScaler是专门为分布式场景设计的它通过一次数据扫描就能计算出每个维度的均值和方差然后应用到每一行数据上from pyspark.ml.feature import StandardScaler scaler StandardScaler(inputColfeatures, outputColscaled_features) scaler_model scaler.fit(train_data_vector) train_data_scaled scaler_model.transform(train_data_vector)注意fit和transform是两个阶段fit是在训练集上计算统计量transform是把这些统计量应用到训练集或者测试集上。在测试集上做标准化时必须使用训练集上拟合好的scaler模型而不能在测试集上重新fit否则会引入数据泄露导致评估结果虚高。这一点和单机版的Scikit-learn的用法完全一致。再就是采样偏差的问题。分布式数据通常量很大很多人为了快速试验会先抽样一部分数据来做特征工程和模型训练。这时候要特别注意抽样方式。如果正负样本比例极不均衡直接用sample(0.1)这种随机抽样会进一步加剧样本不均衡的问题。建议使用分层抽样train_data_sample train_data.stat.sampleBy( label, fractions{0: 0.2, 1: 1.0}, seed42 )这样可以在保持正样本全量参与的情况下对负样本进行下采样既能减少数据量又不会让小类别的样本信息完全丢失。4. 核心算法实战分类、聚类、推荐三大场景的MLlib应用4.1 分类与回归逻辑回归和随机森林是主力MLlib中的分类和回归算法覆盖了大部分工业场景。逻辑回归是二分类任务的首选它的训练过程经过优化在处理大规模稀疏特征时非常高效。使用起来很简洁from pyspark.ml.classification import LogisticRegression lr LogisticRegression(featuresColscaled_features, labelCollabel) lr_model lr.fit(train_data_scaled)在训练之前可以通过setParams或者直接在构造时传入超参数比如正则化系数regParam、弹性网络混合比elasticNetParam。逻辑回归的正则化非常重要尤其是高维稀疏特征的场景下不加正则化基本会过拟合。实际调参的时候我一般会把regParam搜索范围设置在0.001到0.1之间。随机森林则是我在处理非线性关系时的首选。它不需要做特征标准化对缺失值也有天然的容忍度而且可以直接输出特征重要性帮助我们做特征筛选from pyspark.ml.classification import RandomForestClassifier rf RandomForestClassifier( featuresColscaled_features, labelCollabel, numTrees100, maxDepth10, seed42 ) rf_model rf.fit(train_data_scaled) importance rf_model.featureImportances关于maxDepth的取值有一个原则数据量越大、特征越多树的深度和数量通常也需要相应增加但要注意训练时长和过拟合风险。在分布式环境下numTrees越大各个树的训练并行度越高性能会更好。我实际跑下来100棵树和10层的深度在大多数业务场景里是一个性价比很高的起点可以先跑出一个基线再通过交叉验证继续调优。这里提一下交叉验证的并行度问题。MLlib里的CrossValidator支持并行评估不同的参数组合但要注意设置setParallelism参数并且理解它的瓶颈所在。默认情况下CrossValidator是基于参数组合做并行化的而不是基于数据或者fold做并行化。如果你的ParamGridBuilder只搜索了一两个参数并行度可能上不去K折交叉验证的各个fold会串行执行。想要加速需要同时配置spark.sql.shuffle.partitions并且设计足够多的参数组合来充分利用集群资源。4.2 聚类KMeans在用户分群中的实战聚类的典型应用是用户分群、异常检测和向量量化。MLlib中的KMeans非常适合处理海量用户特征的分群需求。我自己常用它做用户行为标签的初筛给后续做个性化推荐提供候选集。from pyspark.ml.clustering import KMeans kmeans KMeans(featuresColscaled_features, k10, seed42) kmeans_model kmeans.fit(train_data_scaled) predictions kmeans_model.transform(train_data_scaled)选择K值一直是聚类的一个重点。MLlib提供了计算轮廓系数的内置方法ClusteringEvaluator你可以对不同的K值分别计算轮廓系数选择得分最高或者曲线拐点处的K值from pyspark.ml.evaluation import ClusteringEvaluator evaluator ClusteringEvaluator(featuresColscaled_features, metricNamesilhouette) silhouette_score evaluator.evaluate(predictions)不过做分群时要留意轮廓系数高不代表分群结果对业务有意义。聚类结果最终是要服务于业务动作的你得看每个群在业务指标上是否有明显差异。例如用户分群后A群的付费金额是B群的10倍A群最近一次登录距今天数是1天B群是60天这种差异才说明分群有实际价值。4.3 协同过滤基于ALS的推荐召回推荐系统是MLlib另一大核心应用场景。基于ALSAlternating Least Squares交替最小二乘的协同过滤算法在处理大规模用户-物品交互矩阵时有非常好的表现。它的训练代码很简单from pyspark.ml.recommendation import ALS als ALS( userColuser_id, itemColitem_id, ratingColrating, rank20, maxIter15, regParam0.01, coldStartStrategydrop, seed42 ) als_model als.fit(interactions)rank表示隐向量的维度一般取10到50之间。regParam是正则化参数防止过拟合。maxIter是最大迭代次数ALS一般15到20次就足够收敛。coldStartStrategydrop这个参数非常重要。线上往往会出现训练集中没有出现过的新用户或者新物品。如果不设置这个策略预测时会出现NaN值导致下游任务出错。设置为drop后模型在预测时会自动丢弃这类样本或者在transform结果中不包含这些无法预测的行。ALS还支持显式反馈和隐式反馈两种模式。如果数据是点击、浏览这种隐式反馈需要设置implicitPrefsTrue并且用置信度权重来区分不同强度的交互。比如用户浏览过和用户购买过对模型的置信度是完全不同的。5. 模型训练的资源和性能调优真实集群上的实测记录5.1 作业提交参数别让你的集群闲着这一部分我拿一次实际的模型训练任务来说明。当时的数据量大约4亿行特征维度300个需要训练一个随机森林模型和一个逻辑回归模型。集群规模是20个节点每个节点有16核CPU和64GB内存。第一次提交作业时我用了最简单的参数结果发现训练耗时极长Spark UI的核心利用率不到10%。后来检查才发现就是因为没有设置--executor-cores导致每个Executor只分配到了1个核心。正确提交方式是这样的spark-submit \ --class com.example.ClickThroughRateTrain \ --master yarn \ --deploy-mode cluster \ --num-executors 40 \ --executor-cores 4 \ --executor-memory 12g \ --driver-memory 8g \ --conf spark.sql.shuffle.partitions200 \ --conf spark.memory.fraction0.6 \ --conf spark.memory.storageFraction0.5 \ click_through_rate.jar参数的含义解释一下--num-executors 40总共启动40个Executor。--executor-cores 4每个Executor使用4个CPU核心。--executor-memory 12g每个Executor分配12GB内存。--driver-memory 8gDriver端分配8GB内存用于收集结果和控制调度。40个Executor乘以4核心总共有160个核心可以并行计算比之前的40个核心效率提升了3到4倍。数据量在4亿行时随机森林训练时间从3小时降到了45分钟。5.2 内存模型搞懂Executor内存就搞懂了OOMSpark的内存模型是整个性能调优里最核心也最容易被忽视的部分。从Spark 2.x开始Executor的内存被统一管理分为执行内存Execution Memory和存储内存Storage Memory两部分它们共享同一个内存区域通过spark.memory.fraction控制占比。spark.memory.fraction默认是0.6表示Executor内存中用于Spark内部执行和缓存的比例剩下的0.4留给用户代码和系统开销。spark.memory.storageFraction默认是0.5表示在共享区域内存储内存占初始比例的一半。训练模型时如果数据量特别大Shuffle阶段需要大量的执行内存。这时候OOM的通常不是Driver而是Executor。我的排查经验是先看Spark UI中Executor的GC时间如果GC时间占比超过20%说明Executor内存不足。这时优先考虑增大spark.memory.fraction或者增大--executor-memory。但是如果Executor本身分配的内存就超过物理机内存YARN会直接杀掉Container这时反而要考虑减少单Executor内存、增加Executor数量。有一个经验公式可以帮助你估算合理的Executor内存$$executorMemory nodeMemory \times 0.75 / executorCountPerNode$$比如节点是64GB内存预留25%给系统和其他组件剩下48GB分给Executor。如果单个节点只跑1个Executor可以给48GB如果跑2个Executor每个给24GB。这样既充分利用内存又避免Container因超内存被杀。5.3 常见性能瓶颈的定位思路与调优记录实际训练中我遇到过的性能瓶颈主要集中在以下三个位置第一个是Shuffle阶段。当特征工程涉及大量的groupBy和join时Shuffle的数据量非常大。定位方法是在Spark UI的SQL Tab里查看各个Stage的Shuffle Read/Write大小。如果单个Stage的Shuffle Write超过几GB就要考虑调整spark.sql.shuffle.partitions。我在4亿行数据做特征拼接时把这个参数从默认的200调到了800单个Stage的处理时间从15分钟降到了6分钟。第二个是数据倾斜。表现在某些Task处理的数据量远大于其他Task整个Stage的完成时间被少数几个慢任务拖死。常见的解决办法是加盐salting或者做广播Join。对于MLlib训练阶段如果label列极度不均衡比如正样本只有0.1%可以考虑在训练前的抽样阶段做分层处理而不是把不均衡数据直接喂给算法。第三个是Driver端OOM。当模型要对全量数据进行某种聚合后收集到Driver端时Driver内存很容易爆掉。比如决策树模型训练完成后通过featureImportances打印不涉及收集全量数据但在做collect操作时就非常危险。处理这类问题时要么增加Driver内存要么尽量避免大结果集的collect操作。6. Pipeline与模型落地从训练到线上预测的完整闭环6.1 Pipeline机制把特征工程和模型训练串起来MLlib借鉴了Scikit-learn的Pipeline思想用Transformer和Estimator构建机器学习工作流。这在实际项目里的价值非常大你的特征处理流程和模型训练可以在同一个Pipeline里定义之后在训练集上fit在测试集或线上数据上transform不会出现训练与预测时特征处理不一致的问题。一个典型的Pipeline示例from pyspark.ml import Pipeline from pyspark.ml.feature import StringIndexer, VectorAssembler, StandardScaler from pyspark.ml.classification import RandomForestClassifier string_indexer StringIndexer(inputColuser_id, outputColuser_index, handleInvalidkeep) assembler VectorAssembler(inputCols[user_index, age, order_cnt, total_amount], outputColraw_features) scaler StandardScaler(inputColraw_features, outputColfeatures) rf RandomForestClassifier(featuresColfeatures, labelCollabel, numTrees100) pipeline Pipeline(stages[string_indexer, assembler, scaler, rf]) pipeline_model pipeline.fit(train_data)这样做的好处是线上预测时只需要加载这个pipeline_model然后直接调用transform方法Spark会自动执行所有的特征处理步骤不需要你在线上重新编写一套特征处理逻辑。这极大地降低了训练和预测之间的不一致风险。6.2 模型保存与加载可靠落地的关键PipelineModel本身可以很方便地保存到分布式存储中pipeline_model.write().overwrite().save(hdfs://namenode:8020/models/user_churn_model)加载时只需要一行代码from pyspark.ml import PipelineModel loaded_model PipelineModel.load(hdfs://namenode:8020/models/user_churn_model)这里有几个要注意的坑。第一模型保存的路径最好是带版本号的比如user_churn_model_v1.0.3便于回滚和对比。第二保存前要确认模型是否还在YARN集群上运行——如果提交的是--deploy-mode client模型默认保存在Driver节点的本地文件系统上只有指定了HDFS路径才会保存到分布式文件系统里。第三加载模型时依赖的JAR包版本需要和训练时一致否则在反序列化阶段可能会报ClassNotFoundException或NoSuchMethodError。6.3 线上预测的两种常见架构模型训练好之后剩下的问题是线上怎么用。我总结下来主要有两种落地方案。第一种是离线批量预测。适合场景是每天凌晨对全量用户跑一次预测把结果写入HBase或者数据库供业务查询使用。这种方案非常简单直接用Spark程序读取前一天的数据加载模型批量transform然后把结果写出去。成本低维护简单适合对实时性不高的场景。第二种是实时在线预测。适合场景是对实时性要求高的场景比如用户在页面上的实时推荐。这种方案一般有两条路一条是用Spark Structured Streaming实时读取Kafka流数据在流上做特征工程和预测另一条是把模型导出为PMML格式部署到Java服务或者规则引擎里做毫秒级预测。关于后者MLlib对PMML导出的支持程度有限目前主要支持的还是传统的回归和聚类模型。如果你用的是树模型或者集成模型导出PMML会有一定的兼容性风险。相比之下我更推荐Streaming方案因为Spark的流处理和批处理可以复用同一套Pipeline代码运维心智负担更小。7. 通过几个真实的坑聊聊MLlib调试的经验7.1 坑一StringIndexer在处理新类别时的静默失败这个坑是在一次上线事故里暴露出来的。训练集里用户ID编码成了1到100万模型上线后发现部分用户的行为特征里出现了异常值进而在预测时全部给出了极端的风险分数。排查发现原因很简单测试集里出现了训练集中没见过的用户IDStringIndexer默认对未见过的类别设置了error策略也就是直接抛异常但因为我们提前设了handleInvalidkeep这些新ID被统一放到了最后一个桶里反而造成了数值上的偏差。正确的做法是在StringIndexer中显式设置handleInvalidkeep然后配合Imputer或Bucketizer对缺失值做单独的映射或者干脆把这类新类别处理为特殊的语义代表比如“未知用户”。更稳妥的方案是对用户ID这类高基数的类别特征不做StringIndexer而是用哈希分桶或者CountEncoder这类方式避免类别数膨胀带来的风险。7.2 坑二AUC虚高——做特征工程时把未来信息泄露进去了这是我见过的最隐蔽的问题。在流失预测项目中我一开始构造的特征里包含了一个“近30天消费金额”的字段。表面上看很合理但在做训练集和测试集切分时我没有严格按照时间顺序切分而是随机切分。这导致同一时间段的数据同时出现在训练集和测试集里模型学到了时间窗口内的数据规律在测试集上的AUC高达0.97。后来按时间序列重新切分数据后AUC直接回落到0.81。这个坑告诉我们一个原则凡是涉及时间窗口的特征训练集和测试集必须按时间切分严禁随机切分。这个教训在我后续的项目中反复出现每次做特征工程都要先想清楚特征的生成时间确保特征生成时间严格早于标签时间。7.3 坑三模型在训练集上的稳定性不等于大规模数据上的稳定性最后想说一个性能和稳定性相关的问题。在大规模数据上训练时模型可能因为个别异常特征值而震荡。比如某天数据源出现异常某个特征列的值突然变成了一个极大的数StandardScaler算出来的均值和方差会被这个值严重带偏导致整个模型的预测产生系统性偏差。应对这个问题的办法是在特征工程阶段加一个数据质量检测每一列都要做最大值、最小值、空值比例的统计设置合理的告警阈值。MLlib中没有内置这个能力但用Spark SQL做一次全量统计的成本很低几分钟就能跑完却能避免后面更大的麻烦。我在项目中形成了一套很固定的特征检测流程data_summary train_data.select([ F.count(F.when(F.col(c).isNull(), c)).alias(c _null_cnt), F.min(c).alias(c _min), F.max(c).alias(c _max) ] for c in feature_cols]) data_summary.show(verticalTrue)如果发现某些特征的空值比例超过20%或者最大最小值明显超出业务常识就主动干预特征定义而不是让异常值悄悄进入模型训练。最后分享两点我现在养成的习惯如果让我对这篇文章做一个实质性的收尾我想总结两条在工作方式上的体会。第一先把Pipeline建好再调参。我刚用MLlib时习惯先单独跑算法再回头做特征工程导致代码散落各处调试效率很低。现在我会先把整个Pipeline框架搭好哪怕先用很少的数据量跑通然后再用CrossValidator或者TrainValidationSplit调参。这样做的好处是特征工程和算法之间有了一个明确的接口——features列后续不管是换模型还是调特征都不需要改动大框架。第二理解Spark调度比背API重要。真正的大数据AI应用难点往往不在机器学习算法本身而在分布式计算的资源管理和数据组织。你什么时候理解了Executor内存结构理解了Shuffle过程的代价理解了数据倾斜和分区数之间的关系MLlib的很多使用困惑就自然解开了。框架在迭代API在演进但分布式计算的基本原理解释了绝大部分线上性能问题。