ARTICLE DETAIL

资讯详情

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

Spark MLlib特征选择实战:从卡方筛选到树模型重要性

Spark MLlib特征选择实战:从卡方筛选到树模型重要性 1. 项目概述与整体设计思路1.1 为什么在大数据场景下做特征选择做算法开发的朋友一定知道特征选择不是锦上添花而是直接影响模型上线的关键一环。在Apache Spark这种分布式计算框架里跑特征选择跟单机环境有本质区别数据量动辄千万级甚至亿级特征维度往往成百上千如果一股脑把所有特征都喂给模型训练时间翻倍不说还容易引入噪声导致模型泛化能力变差。我在实际项目中遇到过一个典型场景某电商平台的用户行为数据原始特征经过清洗、拼接、窗口统计之后维度膨胀到了1200多个。直接丢进Spark训练的GBDT模型里跑一轮集群资源占用率飙到90%以上训练耗时从40分钟涨到3小时。后来在特征选择阶段下功夫把维度压到180个训练时间回到50分钟以内AUC还从0.72提升到了0.78。这个案例给我的教训是特征选择做得好模型效果和资源开销是双赢的。1.2 特征选择的三个层次特征选择不只是删掉几个没用的列这么简单。我习惯把它拆成三个层次来看第一层是粗糙过滤。先看特征的基本统计量缺失率超过80%的、方差几乎为0的、跟目标值完全不相关的直接剔除。这一步成本极低在Spark里做基本的统计聚合就行。第二层是算法筛选。用卡方检验、F检验、互信息这类统计方法给每个特征打个分按得分排序保留Top N。这个层次已经能让很多无关特征现出原形。第三层是包裹式/嵌入式选择。这是最贵但也更可靠的做法——用模型自身的反馈来判断特征的重要性比如逻辑回归的系数绝对值、随机森林的基尼重要性、L1正则化自动把系数压缩到0等。Spark MLlib对这三层都有对应的实现。很多初学者容易犯一个错误一上来就直接用最高级的方法忽略了分层筛选的价值。我的经验是先快后稳先用粗过滤把明显没用的删掉再用中等成本的方法做一轮最后结合模型反馈做精细筛选。这样整个流程可控、可解释、也更容易排查问题。1.3 一个完整的Spark特征选择闭环我们在做特征选择的时候必须想清楚一件事选择出的特征最终要服务于后续的模型训练和推理。所以我不建议把特征选择当成一个孤立的步骤而是应该把它放进整个Pipeline里统一管理。在Spark里Pipeline机制可以很方便地把特征处理、特征选择、模型训练串联起来。我用到的链路大概是这样的数据清洗 → 缺失值填充 → 类别编码StringIndexer / OneHotEncoder → 特征拼接VectorAssembler → 特征选择ChiSqSelector / 模型重要性 → 训练模型 → 评估与调参。关键是特征选择的结果要作为Pipeline的一部分保存下来这样在推理阶段遇到新数据时可以通过同一个Pipeline做完全一致的特征变换和特征筛选避免训练时用的特征和线上推理时不一致的经典翻车事故。1.4 本项目适用的读者与场景这篇指导主要面向三类读者第一类是刚接触Spark算法开发、想了解特征选择怎么落地的初中级工程师第二类是在做大规模特征工程、希望通过特征选择降低训练成本的算法工程师第三类是需要在集群环境里跑离线训练的机器学习平台开发人员。不管你用的是Spark MLlib自带的特征选择API还是用Spark做数据预处理后把结果导出给XGBoost、LightGBM等框架使用本文的思路都适用。我下面分享的内容是基于2.4版本以上的Spark MLlib部分API在2.4以下会有差异注意区分。2. Spark MLlib特征选择核心工具解析2.1 ChiSqSelector卡方选择器的原理与适用边界Spark MLlib里最常用的特征选择器之一是ChiSqSelector也就是卡方特征选择器。它的核心思想很简单对每个特征与标签做独立性检验卡方值越大说明特征与标签的关联性越强反之则关联性弱。这个方法的计算逻辑在统计学教材里写得很清楚先假设特征与标签无关然后计算实际观测频数与期望频数之间的偏离程度偏离越大卡方统计量越大越倾向于拒绝无关的原假设。在分布式环境里Spark会把你提供的特征和标签做交叉汇总计算每个特征每个取值下的卡方值然后按重要性排序。不过我要提醒一点卡方检验只适合类别型特征与类别型标签。对于连续型数值特征输入到ChiSqSelector之前通常是先做离散化比如分箱或者干脆不用卡方选择器。我踩过这样的坑数据里有个年龄字段取值范围在18到65之间直接把它以连续值的形态丢给ChiSqSelector结果跑了半天还报错提示要求非负离散值。所以在做卡方选择之前必须确认特征的数据类型和取值形态。ChiSqSelector提供的参数包括numTopFeatures保留Top N个特征、percentile按比例保留、fpr通过控制假阳性率来选择、fdr控制错误发现率、fwe控制族错误率。我用得最多的是numTopFeatures和percentile因为这两者的业务语义最直观其他几个方法在样本量特别大时会偏向选择更多特征反而不容易调。2.2 VectorSlicer手动精确保留指定维度如果说ChiSqSelector是机器帮你选那VectorSlicer就是你自己拍板。它的作用是从已有的特征向量里直接切出指定的列保留你关心的维度。这个工具特别适合在你不完全信任自动特征选择、但已经依靠业务经验明确知道哪些特征必须保留的时候使用。比如信贷风控项目里身份证归属地、近3个月查询次数、公积金缴纳比例这些特征是业务强相关的即便统计得分一般也要手动留下。这时候VectorSlicer就派上用场了你可以按索引切也可以按列名切。需要注意的是VectorSlicer的输入必须是向量类型的列所以通常放在VectorAssembler之后使用。它并不做任何判断只是按位置索引或名称做裁剪。在Pipeline里它的位置通常在最前面像一道闸门一样控制着哪些特征能进入后续的自动筛选流程。2.3 PCA主成分分析的降维本质与取舍PCA在特征工程里通常被归为降维方法但它其实也承担了一部分特征提取和特征筛选的功能。PCA的核心思路是把原始高维空间映射到一个低维子空间让映射后的各主成分按方差贡献度从高到低排列前几个主成分往往集中了绝大部分信息。我在实际中用PCA的场景有两种一种是特征维度实在太多而且特征之间相关性很强例如同样的用户行为统计在多个时间窗口内重复出现直接做PCA能大幅压缩维度另一种是后续模型对特征独立性有要求比如部分线性模型PCA可以有效降低共线性带来的不稳定风险。但PCA的缺点也很明显主成分是原有特征的线性组合可解释性极差。你在给业务方解释模型的时候很难说清楚第三主成分到底代表什么业务含义。所以我的原则是如果业务解释性要求高尽量不用PCA如果只是为了压缩维度、追求模型指标那PCA是个高效的选择。Spark MLlib里PCA的使用很简单主要参数是k要保留的主成分数量。我一般会结合累计方差解释率来选择k比如累计解释率达到85%以上就认为信息保留得差不多了。2.4 基于模型的特征重要性筛选除了上面几个专用工具Spark里还有一种非常灵活的特征选择方式——借助树模型的特征重要性。RandomForestClassifier训练完之后可以通过featureImportances属性拿到每个特征的相对重要度排序之后截取Top N。这个方法的原理是树模型在分裂时会反复利用区分度高的特征如果某个特征在节点分裂中贡献的基尼杂质降低越多它的重要性就越高。拿到的featureImportances是一个稀疏向量非零元素对应参与训练的特征不包括被抛掉的排序后就能定位重要特征。我通常在两个阶段使用这个属性第一个阶段是探索性分析先快速跑一个浅层随机森林看特征重要性的排序大致是什么样子帮助我发现异常特征比如一个本不该重要的特征排到了第一这往往暗示数据泄漏第二个阶段是最终筛特征用较深的随机森林跑稳定后把重要性低于阈值的特征删掉。有一个细节容易踩坑RandomForest的featureImportances受到特征量纲影响很大。如果某个特征是0到1的比例值另一个特征是几万量级的金额值模型分裂时对量纲大的特征更敏感重要性容易被高估。因此在基于模型做特征选择前最好先用VectorAssembler拼接特征并在之前对连续特征做标准化或归一化。2.5 各方法对比与选型建议我发现很多同学在选特征选择方法时拿不定主意这里给一个简单的对比表结合我的实践经验帮助大家快速做出选择方法适用特征类型标签类型可解释性计算成本适用场景ChiSqSelector类别特征离散类别标签高低分类问题初筛VectorSlicer任意任意高极低业务指定特征保留PCA连续数值特征无关低中高维且共线性严重RandomForest重要性数值/类别均可类别/回归中较高精细筛选与探索分析选型逻辑我一般这样走先问自己删特征后要不要给业务解释要的话优先考虑VectorSlicer和ChiSqSelector接着看特征类型类别型特征多就多依赖卡方连续型特征多就多依赖特征重要性排序最后看数据量和算力预算预算紧张就少跑几轮复杂的树模型。3. 实操过程与核心环节实现3.1 实验环境与数据集说明为了让大家能直接复现我用一套公开可模拟的数据来做演示数据集特征是混合型的包含数值特征和类别特征标签是二分类。数据规模我故意设置得大一点模拟真实业务场景的千万行级别这样才能体现出Spark分布式计算的优势。我这里用Spark MLlib的Scala代码来做演示Python版API逻辑完全一致注意参数名称对应即可。下面是我的基础环境配置import org.apache.spark.sql.SparkSession val spark SparkSession.builder() .appName(FeatureSelectionExample) .master(yarn) .config(spark.sql.shuffle.partitions, 200) .config(spark.executor.memory, 8g) .getOrCreate()shuffle分区设置成200是经验值。数据量在千万行级别、单分区几万行左右时200个分区调度开销小、并发度也够用。太小了容易出现数据倾斜太大了调度器本身会有压力。如果你跑到亿级数据建议把分区调到500到800。3.2 数据清洗与特征向量构建拿到原始数据之后第一步一定是清洗和构造特征向量。这里的顺序很有讲究先处理缺失值再进行类别编码最后拼装成向量。缺失值处理放在前面是为了避免后面编码和拼接的过程中把缺失值当成特殊类别导致特征分布扭曲。import org.apache.spark.ml.feature._ import org.apache.spark.ml.Pipeline // 假设原始数据包含age, income, city, gender, is_click标签 val df spark.read.parquet(/data/user_features/) // 1. 缺失值填充数值列用中位数类别列用模式值 val imputer new Imputer() .setInputCols(Array(age, income)) .setOutputCols(Array(age_imp, income_imp)) .setStrategy(median) // 2. 类别特征编码 val genderIndexer new StringIndexer() .setInputCol(gender) .setOutputCol(gender_index) val cityIndexer new StringIndexer() .setInputCol(city) .setOutputCol(city_index) // 3. 对类别索引做OneHot防止模型误把类别ID当成有大小关系的数值 val genderEncoder new OneHotEncoder() .setInputCol(gender_index) .setOutputCol(gender_vec) val cityEncoder new OneHotEncoder() .setInputCol(city_index) .setOutputCol(city_vec) // 4. 拼接特征向量 val assembler new VectorAssembler() .setInputCols(Array(age_imp, income_imp, gender_vec, city_vec)) .setOutputCol(raw_features)StringIndexer对类别列做编号之后如果直接用这个编号当特征喂给模型类别之间就会被赋予人为的排序关系这完全不合理。OneHotEncoder则把类别索引展开成多个0/1维度规避了此问题。这里有个实际经验想说OneHot编码之后特征维度会膨胀尤其城市这种高基数类别特征。如果你不做特征选择直接拼一个几千维的稀疏向量去训练虽然SKlearn之类的单机库能扛但在Spark集群里完全是浪费资源。所以我在OneHot这一步之后通常紧接着做一轮卡方筛选。3.3 用ChiSqSelector筛选Top特征将原始特征向量构建完成后就可以调用ChiSqSelector准备卡方筛选了。但直接拿连续型数值特征进卡方选择器是会出问题的正如前面提到的那样卡方要求的输入是非负离散值。所以我在卡方选择之前把连续字段做了分箱处理import org.apache.spark.ml.feature.{Bucketizer, ChiSqSelector} // 对age和income做分箱处理 val ageSplits Array(Double.NegativeInfinity, 20, 30, 40, 50, 60, Double.PositiveInfinity) val bucketizer new Bucketizer() .setInputCol(age_imp) .setOutputCol(age_bucket) .setSplits(ageSplits) // 卡方选择器保留Top 30个特征 val selector new ChiSqSelector() .setFeaturesCol(raw_features) .setLabelCol(is_click) .setOutputCol(selected_features) .setNumTopFeatures(30)分箱边界是根据业务含义设定的年龄按20、30、40、50、60划分这是互联网业务里的常见做法。这里要说明分箱只是为了让连续特征能进入卡方选择流程选定特征之后模型训练阶段仍然可以用原始的连续特征不必纠结于分箱是否损失了精度。对于高基数类别特征比如城市OneHot之后每个取值对应一个维度卡方选择器会在这个稀疏向量上工作。理论上每个类别维度的卡方值单独计算对于出现频次特别低的类别卡方值可能虚高。所以我在实际项目中会先对低频类别做合并比如出现次数低于总样本量万分之一的统一归为其他这能有效压制卡方检验的小样本偏误。3.4 通过随机森林特征重要性做二次筛选卡方筛选出来的特征仍然是基于单变量统计的没有考虑特征之间的交互。我习惯在卡方筛选之后再跑一个随机森林用特征重要性做二次排列和筛选。import org.apache.spark.ml.classification.RandomForestClassifier import org.apache.spark.ml.evaluation.BinaryClassificationEvaluator val rf new RandomForestClassifier() .setFeaturesCol(selected_features) .setLabelCol(is_click) .setNumTrees(200) .setMaxDepth(8) val rfModel rf.fit(trainData) // 打印特征重要性 val importance rfModel.featureImportances println(importance.toArray.zipWithIndex .sortBy(-_._1) .take(10) .map { case (score, idx) s特征索引[$idx] 重要性[$score] } .mkString(\n))跑出来的重要性结果我建议不要只看排名还要看分布的陡峭程度。如果大多数特征的重要性都集中在少数几个特征上并且头部特征与剩余特征之间的分数差出一个数量级那说明你只需要保留头部少数特征。如果重要性的下降曲线很平缓说明特征携带的信息比较分散你的截断阈值就可以放宽一些。我通常会做一个累计重要性曲线来辅助决策达到总重要性85%的特征数量就是我的候选保留个数。有了这个数字作为参考回到卡方选择器里调整NumTopFeatures就很有依据不再靠猜。3.5 构建完整Pipeline并对比效果为了演示特征选择的实际收益我把完整的流程封装进Pipeline统一管理。这里有个关键设计特征选择器选出来的特征列直接作为下一阶段分类器的输入特征列。import org.apache.spark.ml.Pipeline val pipeline new Pipeline().setStages(Array( imputer, genderIndexer, genderEncoder, cityIndexer, cityEncoder, bucketizer, assembler, selector, rf )) val model pipeline.fit(trainData)预测时直接调用transform即可val predictions model.transform(testData) predictions.select(is_click, prediction, probability).show()用同样的数据跑了两组对比实验一组是跳过特征选择的完整Pipeline另一组是带特征选择的Pipeline。在相同的随机森林参数下前者的AUC是0.74训练耗时42分钟后者的AUC是0.76训练耗时22分钟。测试集上的结果说明经过选择后的特征组合更紧凑减少了部分噪声特征对树模型的干扰同时训练效率显著提升。当然单次实验有偶然性我后来在几个不同数据集上都做了类似验证结论大体一致在维度较高、噪声较多的场景下特征选择带来的收益是稳定的。维度本身不高的时候特征选择对效果的提升不明显但训练耗时仍然会有可观的下降。3.6 保存Pipeline用于线上推理模型训练完成后把Pipeline保存下来线上推理时直接加载使用。这一步我必须强调因为很多人会忘记保存Pipeline和保存模型不是一回事Pipeline里包含了所有特征变换和特征选择的逻辑只存模型的话推理时还得重新走一遍特征处理工作量大且容易出错。pipelineModel.write.overwrite().save(hdfs:///models/feature_select_rf)线上服务读取这个目录拿新数据一transform就能得到预测结果。所有分箱边界、编码映射、特征选中的索引都已经固化在Pipeline里不会和训练阶段产生偏差。4. 常见问题与排查技巧实录4.1 卡方选择器报Negative values错误这是我被问到最多的问题。ChiSqSelector对输入的非负性有严格要求连续型特征如果有负值比如标准化之后的Z-score特征直接报错并且错误信息还不算友好。这个问题的根源在于卡方检验的本质它基于频数表计算观察频数和期望频数频数不可能是负数。解决方案有两个一是对连续特征做分箱用分箱索引代替原始数值二是改用其他不以频数为基础的特征选择方法比如ANOVA方差分析在Spark里通过ANOVASelector提供它就能处理连续特征。需要提醒的是不要试图靠给负值加一个偏移量比如全部加100来骗过卡方检验这种做法会严重影响检验结果毫无统计学依据。4.2 StringIndexer出现无法处理的标签StringIndexer默认遇到模型未见过的类别会抛异常这在线上推理时尤其坑。更常见的情况是训练数据里有些类别出现频次极低StringIndexer依然会给它分配索引但这些类别在测试集里几乎不出现导致稀疏维度。我的处理方法是在使用StringIndexer之前先对低频类别做合并。Spark没有现成的按频次合并API需要自己写一段聚合处理也可以用QuantileDiscretizer做类似归类。做完之后再走编码流程模型就稳定很多。4.3 数据倾斜导致训练任务卡死特征选择阶段经常需要做聚合操作当数据里某个类别的占比特别高比如二分类任务里正样本只占1%聚合计算就容易出现数据倾斜。表现是某个Executor跑得非常慢其余Executor都在空等整个任务卡在那里。排查技巧是看Spark UI里的Stage执行时间分布如果单个任务比其他任务慢出一个数量级基本可以判定是倾斜。常见的解法有三个一是增加shuffle分区数让倾斜的Key分散到更多分区二是对热点Key做局部加盐先局部聚合再整体聚合三是在做统计计算时用approxQuantile或近似Count-Min Sketch代替精确计算。特征选择这个场景里我推荐用近似方法因为我们要的是特征重要性排序不是精确的统计量近似结果完全够用但性能提升是十倍级别的。4.4 稀疏特征向量与稠密向量混用报错在Pipeline里特征向量可能是SparseVector也可能是DenseVector这取决于上游的构造方式。OneHotEncoder输出的是稀疏向量VectorAssembler拼接时如果只有稀疏向量参与输出也是稀疏的。但如果你在后面插入了一些产生稠密向量的操作比如StandardScaler输出稠密再往后续环节传类型不一致就会报错。我的建议是整个特征处理链路上尽量保持特征向量的稀疏性。具体做法是在VectorAssembler之前明确控制哪些Transform输出的向量类型并在必要的时候用Vector的toSparse方法做显式转换。否则一旦维度膨胀稠密向量会带来巨大的内存开销执行速度也会明显下降。4.5 特征选择过拟合问题特征选择本身也可能过拟合这一点经常被忽略。尤其当你对同一份数据反复跑多轮特征选择、每次都根据结果微调阈值最后选出的特征会越来越贴合训练集在测试集上效果却并不理想。我的做法是把特征选择的过程放在交叉验证的外层也就是先选特征再对选中的特征做交叉验证这样能避免在测试集上调特征选择参数的隐性泄漏。如果数据量足够也可以把数据集拆成三份一份用于特征选择、一份用于训练、一份用于最终验证。4.6 特征维度选择阈值没有统一标准有些同学总想找到一个标准答案——到底保留多少个特征最合适这个问题没有通用答案但有一个朴素的经验法则从业务角度先划定一个候选区间比如30到200然后在这个区间里每隔20个特征取一个点跑一组模型画出特征数量-模型指标曲线取曲线拐点的位置作为阈值。我见过很多项目在这个拐点上被打脸原因是只看了AUC没有考虑训练耗时和特征维护成本。实际工作中如果增加20个特征只带来0.001的AUC提升但每月数据维护成本成倍上升那这20个特征就应该被砍掉。特征选择最终要服务于工程和业务的综合利益不只是模型指标。5. 实操心得与扩展建议5.1 我踩过的最隐蔽的一个坑最后分享一个让我印象特别深的教训某次特征选择做完模型指标看起来很不错但我发现一个应该是敏感信息的特征排在了重要性第一。排查之后发现这个特征是从某个下游行为表里join进来的而这张行为表本身包含了标签发生之后的信息也就是传说中的数据泄漏。这个坑很隐蔽因为Pipeline和特征选择器都没有任何报错的机制它们只会诚实反映数据的规律。所以大家在拿到特征重要性结果时一定要重审一遍Top特征的业务含义问自己这个特征在预测时点真的能拿到吗答案是不知道或者不确定的哪怕指标再好也要剔掉。特征选择是一个技术过程但最后的验收需要业务判断来兜底。5.2 从Spark到XGBoost/LightGBM的衔接建议很多团队只用Spark来做数据处理和特征选择后续训练交给XGBoost或LightGBM。这种情况下Spark侧特征选择选中的特征索引要导出成一个白名单后续特征平台上的计算逻辑要跟这个白名单保持严格一致。我建议把这件事做成自动化Spark训练完特征选择器后直接把选中的特征名列表写入特征仓库的元数据表下游训练和推理都从元数据表读取而不是各自维护一份。这样一整条链路的变化是批次可追溯的也避免训练用的特征跟特征仓库定义不一致这类问题反复出现。5.3 可以继续扩展的方向如果这篇指导对你有帮助你还可以继续深入几个方向一是用MLlib里的ANOVASelector和Fisher Exact Test做补充尤其是处理连续特征和稀疏二元特征二是结合特征分组做组内选择组间平衡在特征选择中加入业务先验避免某个业务域的特征数量过多挤压其他业务域三是把特征选择纳入自动机器学习框架在超参数搜索的同时调整特征数量让特征维度和模型参数的组合一起优化。我自己现在做新项目时特征选择已经不是临时抱佛脚的环节而是贯穿所有算法开发流程的基础设施部分。希望这篇内容能帮你在Spark算法开发路上少走一些弯路把更多时间花在真正有价值的数据理解和模型优化上。
返回列表