ARTICLE DETAIL

资讯详情

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

Apache Spark特征选择实战:从ChiSqSelector到Pipeline的落地指南

Apache Spark特征选择实战:从ChiSqSelector到Pipeline的落地指南 做算法开发的人大概都遇到过这种场景离线训练集上AUC明明不错换一批线上数据就崩翻来覆去找原因最后定位到建模时喂进去的那几百个特征里有一半以上是凑数的。这个问题放到Apache Spark上会变得更麻烦——特征数量多、数据规模大、Spark又是按Stage分布式计算的多留一个无意义的高基数字段Shuffle和Driver端序列化开销都会成倍上涨。所以特征选择在Spark算法开发流程里不是锦上添花而是决定Pipeline能不能稳定跑下去的必修课。这篇文章不打算写成某个孤立的API手册而是想把这些年把特征选择真正落地到Spark训练Pipeline里的完整思路整理出来哪些方法在分布式环境下实用哪些看起来很美但根本跑不起来每一步工程上有哪些坑选完之后还要做什么。适合正在写Spark ML/MLlib训练代码、被特征维度搞到头大、或者刚接手一个宽表模型的算法开发同学参考。1. 为什么在Spark上做特征选择和单机Sklearn的思路完全不同很多人在单机上用Sklearn的习惯是pandas读进来SelectKBest、RFE、SelectFromModel一顿操作然后拿着选出来的列直接训练。这套思路搬到Spark集群上第一步就会出问题——Spark的核心数据抽象是分布式的DataFrame/RDD数据不会乖乖躺在单机内存里等你numpy运算而特征选择算法本身又很吃数据遍历和迭代次数所以必须重新思考哪些方法在分布式环境下还成立。1.1 数据规模决定算法选型Wrapper在分布式环境下基本不可行特征选择方法按原理分三类Filter过滤式、Wrapper包裹式、Embedded嵌入式。单机上大家惯用的RFE递归特征消除、前向/后向搜索都属于Wrapper它们要反复训练模型、反复用验证集打分特征组合空间稍微大一点就是指数级搜索。在Spark上做这件事等于每一轮特征组合都要触发一次全量数据扫描和模型训练跑一次实验可能就要等几个小时最后选出来的特征组合还不一定比Filter业务规则好多少。所以我个人的结论很直接在Spark上Wrapper基本是一个伪选项除非你的特征量小到可以忽略不计否则不要尝试自己写分布式RFE。真正能在集群上稳定落地的是Filter和Embedded两种路线下面几章会分别展开。1.2 Spark MLlib原生支持的特征选择工具清单很多Spark新手会问MLlib里到底有没有特征选择器有但数量不多而且分布在不同的包下面。我常用的清单大概是这样的方法工具类适用场景特性按索引/名称切片VectorSlicer已有明确特征名单配合VectorAssembler使用卡方检验ChiSqSelector离散特征、分类标签Filter路线只能选不能评估组合树模型重要性RandomForest/GBT 的 featureImportances连续特征、非线性关系Embedded路线效果好L1正则/弹性网络LinearRegression / LogisticRegression 的 elasticNetParam线性场景、特征共线性高Embedded路线系数自然稀疏公式选择RFormula特征少、公式简单便捷但调试困难这里面ChiSqSelector和树模型featureImportances是我最常用的两个。L1正则也值得提一句如果你确认自己的问题是线性可分、特征之间又没有太强的非线性交互直接在Spark ML里把LogisticRegression的elasticNetParam调成1.0训练完看系数是否为0就能顺带完成一次特征筛选成本非常低。1.3 什么时候可以跳过特征选择这一步不是所有项目都需要做严格的特征选择。我见过一些人无论什么模型都把ChiSqSelector挂上去反而把有效特征滤掉了。我的判断标准是如果特征总数只有几十个而且你用的是树模型先不要急着做特征选择树模型的特征重要性自己会处理一部分噪声直接训练后看importance报告就够了。如果你用了深度学习、特征会经过Embedding层传统特征选择的优先级也会下降Embedding本身就是一种端到端的特征表示学习。如果你的模型是线性模型且特征有强共线性优先考虑L1正则而不是卡方因为L1能同时处理共线性和稀疏化。反过来一旦特征数量到了上百甚至上千特征之间又有明显的信息重叠或者你的Pipeline需要经常重训、需要定时做特征监控那特征选择就必须成为固定环节。2. 动手跑算法之前先用Spark SQL把候选特征过滤一轮这里说的不是用算法过滤而是用SQL和业务常识过滤。我见过不少团队直接把宽表里的两百多个字段一股脑塞进VectorAssembler然后再跑ChiSqSelector结果跑出来的重要特征里往往包含了唯一ID、未来指标、或者是已经退市的字段。特征选择的第一步永远是在数据层面剔除垃圾而不是靠算法替你判断。2.1 用数据字典和字段血缘建立特征准入清单动手之前先花半天时间做一件事收集你手上这张特征宽表的数据字典和字段血缘。搞清楚每个字段的生成逻辑、更新频率、口径定义。重点排查三类问题这个字段是否包含了未来信息比如预测用户明天是否退款特征里却有是否已经退款这就是典型的标签泄漏。训练时和推理时这个字段都能拿到吗有些特征在离线ETL里很容易算但线上实时请求时根本拿不到这类字段一旦进入模型上线就会面临特征缺失。这个字段是否已经在其他地方被替代宽表里经常出现同一业务实体的不同粒度汇总比如用户近7天消费金额和用户近30天消费金额它们高度相关没必要全部保留。2.2 缺失率、常量列、超高基数三句SQL解决一半问题业务过滤之后再用Spark SQL做一轮粗筛。以下三句SQL是我在特征探查阶段必跑的-- 常量特征distinct后只有1个值直接踢掉 SELECT col_name, COUNT(DISTINCT col) AS distinct_cnt FROM feature_table GROUP BY col_name; -- 缺失率看每列空值占比 SELECT SUM(CASE WHEN col1 IS NULL THEN 1 ELSE 0 END) / COUNT(*) AS miss_rate_col1, SUM(CASE WHEN col2 IS NULL THEN 1 ELSE 0 END) / COUNT(*) AS miss_rate_col2 FROM feature_table; -- 高基数离散字段比如user_id这种看看基数是不是已经接近行数 SELECT COUNT(DISTINCT user_id) AS user_cardinality, COUNT(*) AS total_rows FROM feature_table;常量特征在高版本Spark里可以直接用VarianceThreshold思路判断但更快的做法还是SQL先数一轮distinct。缺失率超过80%的字段除非业务上明确说缺失本身有含义否则直接剔除。高基数离散字段要特别注意像user_id、order_id这种粒度ID类字段如果直接做StringIndexer再进卡方会产生极其荒谬的排名。2.3 把业务规则沉淀成可复用的特征准入checklist做完了上面的步骤建议不要只是口头把结果同步给同事而是沉淀成一份checklist我自己的项目里通常会长这样是否有未来信息训练和推理口径是否一致线上是否能实时获取延迟和成本是否可接受缺失率是否超过阈值缺失是否需要单独编码成一类基数是否过高是否需要对高基数列做分箱或目标编码是否与已有特征高度相关业务上是否可以合并特征生成任务是否稳定运行源表是否经常Schema变更这份checklist的价值在于当你后面跑出来的重要特征列表里混进一个唯一ID时你能立刻意识到是前面哪一步漏了而不是一头扎进调参里。3. Filter路线落地ChiSqSelector和卡方检验的工程细节Spark里最常见的Filter式特征选择工具就是ChiSqSelector。它底层用的是卡方检验专门用来衡量离散特征和离散标签之间的相关性。训练集上跑一遍fit就能得到一个特征选择器输入原始特征向量输出筛选后的特征向量。3.1 卡方检验到底在算什么O与E的故事卡方检验背后是一个很朴素的思路如果某个特征和标签完全无关那么在这个特征的不同取值下标签的分布应该大致相同。比如用户所在城市这个特征如果和是否点击无关那么北京用户和上海用户的点击率应该差不多如果事实上差很多就说明这个特征和标签有关联。数学上Spark会为每个特征和标签构建一个列联表然后计算卡方统计量χ² Σ(O - E)² / EO是实际观测频数E是期望频数。期望频数按特征和标签独立的理论计算也就是行合计乘以列合计再除以总样本数。卡方值越大说明实际分布和期望分布偏离越大越说明特征和标签不独立也就越值得保留。提示卡方检验天然适合离散特征比如分箱后的用户行为计数对连续数值型特征效果一般。如果你手里全是连续特征要么先用QuantileDiscretizer做分箱再跑卡方要么直接跳到第四章的树模型方法。3.2 ChiSqSelector参数怎么设numTopFeatures、percentile、fpr、fdr、fweChiSqSelector最核心的参数是selectorType有五个选项selectorType含义适用场景numTopFeatures保留卡方值最大的前K个特征业务明确要Top Npercentile保留前百分之多少的特征不知道具体K先粗筛fpr按p值阈值筛选控制单次检验假阳性率特征量大需要统计显著fdr按Benjamini-Hochberg控制错误发现率特征量大且关心误选比例fwe按Bonferroni控制族系错误率特征量很大要求极其保守实际项目里我用得最多的是numTopFeatures和percentile因为业务方永远会问最后留下了多少个特征Top N是最容易沟通的表述。fpr和fdr在特征数量上千、想要统计严谨性时比较有用但要注意它们不控制保留数量只会告诉你哪些显著往往选出来的特征数量很飘忽不好纳入日常Pipeline。代码示例Scala版大概是这样的import org.apache.spark.ml.feature.VectorAssembler import org.apache.spark.ml.feature.ChiSqSelector val assembler new VectorAssembler() .setInputCols(Array(channel_idx, device_idx, hour_idx, city_level, user_star, pay_level)) .setOutputCol(rawFeatures) .setHandleInvalid(keep) val selector new ChiSqSelector() .setFeaturesCol(rawFeatures) .setLabelCol(label) .setOutputCol(selectedFeatures) .setSelectorType(percentile) .setPercentile(0.4) val selectorModel selector.fit(trainDF) val selectedDF selectorModel.transform(trainDF)3.3 一个完整的流水线示例从交易流水里选Top特征假设你有一个交易流水数据集特征是小时区间、设备类型、城市等级、用户星级、支付方式等离散字段。这些字段都是字符串需要先做StringIndexer变成数值索引再统一组装成特征向量import org.apache.spark.ml.feature.StringIndexer import org.apache.spark.ml.Pipeline val indexers Array(channel, device, hour_bucket, city_level) .map { colName new StringIndexer() .setInputCol(colName) .setOutputCol(colName _idx) .setHandleInvalid(keep) } val assembler new VectorAssembler() .setInputCols(indexers.map(_.getOutputCol).toArray) .setOutputCol(rawFeatures) val selector new ChiSqSelector() .setFeaturesCol(rawFeatures) .setLabelCol(label) .setOutputCol(selectedFeatures) .setSelectorType(numTopFeatures) .setNumTopFeatures(20) val pipeline new Pipeline().setStages(indexers : assembler : selector) val pipelineModel pipeline.fit(trainDF) val resultDF pipelineModel.transform(trainDF)这样组织的好处是特征选择直接被封装成Pipeline的一部分后续接分类器也好、保存模型也好都不需要手工维护中间DataFrame。训练完成后pipelineModel里的selectorModel保存了被选特征的索引信息预测时用同一套流程不会出现训练和推理特征顺序不一致的问题。3.4 ChiSqSelector的两个高频翻车点第一个翻车点对StringIndexer后的整数特征做卡方本身没问题但如果你在StringIndexer之前没有处理好缺失值和新类别Spark大概率会抛异常或者在keep模式下把未知类别映射成一个特殊索引测试集预测时和训练集对不上。我的习惯是所有StringIndexer都明确设置setHandleInvalid(keep)并在后续用SQL把它归一化成一个合法的unknown桶。第二个翻车点把连续特征直接丢进ChiSqSelector。很多人以为卡方能处理任意数值实际上是能算算出来的结果却没有太多参考价值。因为连续特征的取值几乎都是唯一的列联表会非常稀疏卡方统计量会被异常值带偏。我摔过一次跟头一个用户累计消费金额特征在卡方排名里排第一替换成同一个人近7天消费金额后模型效果反而更好原因就是累计金额和标签的相关性被极值放大了。连续特征想用卡方先分箱。4. Embedded路线让随机森林的特征重要性替你打分排序如果你的特征以连续值为主、或者特征之间明显存在非线性交互卡方法会显得不太够力。这时候我会换成Embedded路线直接在Spark里训练一个RandomForestClassifier然后读取它的featureImportances。这个方法在训练当前模型的同时完成了特征评估不需要额外的特征选择循环。4.1 featureImportances怎么读归一化相对重要性的含义Spark的树模型特征重要性本质上是对所有节点分裂时带来的不纯度减少量按特征累加最后做归一化所以每个特征的importance加起来等于1。它表达的是这个特征在当前树集合中被用来分裂的贡献占比是一个相对值不是绝对值。import org.apache.spark.ml.classification.RandomForestClassifier val rf new RandomForestClassifier() .setFeaturesCol(rawFeatures) .setLabelCol(label) .setNumTrees(200) .setMaxDepth(10) .setFeatureSubsetStrategy(sqrt) val rfModel rf.fit(trainDF) val featureNames Array(channel_idx, device_idx, hour_idx, city_level, user_star, pay_level) val importancePairs featureNames.zip(rfModel.featureImportances.toArray) .sortBy(-_._2) importancePairs.foreach { case (name, importance) println(s$name: $importance) }读取到importance之后怎么决定保留哪些特征我常用的策略有两种。一是按累计重要性把importance从大到小累加累计到80%或90%时截断剩下的特征数量一般不会太多。二是看排序的肘部画一条重要性曲线曲线在某个位置之后忽然变得平缓那里就是一个自然的截断点。两种方法都比拍脑袋定K值更稳。4.2 用随机森林做评估器还是用GBT我的选择标准Spark里树模型有两个常用选择RandomForest和GBT梯度提升树。它们都能输出featureImportances但我不会随便选。大致判断标准是数据量中等、特征噪声大、希望快速拿到一个稳定的重要性排名选随机森林。它对超参数不那么敏感又容易并行Spark集群上几个Executor就能把200棵树跑完。数据量充足、特征信号更强、愿意多花几倍训练时间换取更尖锐的重要性排序选GBT。GBT在Spark上的训练串行性更强调参也更麻烦maxIter、stepSize、maxDepth都要一起调。如果项目本身最终就要用GBDT类模型上线那直接用GBT算importance就顺带做完了特征选择一举两得。当然importance排名只是给你一个候选名单不一定要严格按Top K截断。我一般会保留累计重要性到85%的特征再结合业务判断补上几个虽然排在后面、但口径明确、线上延迟低的特征。4.3 避免先选特征再训练的过拟合陷阱有个很容易犯的错误先在整个训练集上跑一遍随机森林得到importance然后用筛选后的特征在同一个训练集上再训练一个最终模型然后汇报一个很高的验证集指标。这个指标是虚高的因为特征选择过程已经偷看了验证集里的大量信息。正确的做法应该是要么把特征选择放到交叉验证的Pipeline内部去做让每一折都重新计算importance要么单独划出一个用于特征探索的validation集只在上面做特征选择最终模型在另一部分数据上训练和评估。Pipeline的方式实现起来最干净import org.apache.spark.ml.Pipeline import org.apache.spark.ml.tuning.CrossValidator import org.apache.spark.ml.evaluation.BinaryClassificationEvaluator val stages Array(assembler, rf) val pipeline new Pipeline().setStages(stages) val cv new CrossValidator() .setEstimator(pipeline) .setEvaluator(new BinaryClassificationEvaluator()) .setNumFolds(5) .setParallelism(4) val cvModel cv.fit(trainDF)这样每一折训练出来的rfModel都会输出一组featureImportances交叉验证结束后你可以把所有折的importance平均一下得到更可信的特征排名。代价是训练耗时乘以折数所以特征数量很大的时候我倾向于先在trainDF的子样本上算一轮importance锁定候选名单再做交叉验证。4.4 高基数连续特征是importance最容易骗你的地方树模型的importance有个顽疾它偏爱高基数的连续特征。一个跨越范围很大的特征比如累计消费金额或者用户活跃天数树上会有更多机会拿它做分裂重要性虚高一个取值很少但实际业务含义很强的二值特征反而容易被淹没。判断这个问题的办法有两个把重要性排名和业务经验对照如果一个业务上公认没用的高基数连续特征排名特别高就要警惕。尝试做Permutation Importance也就是把特征值随机打乱后再看模型指标下降多少。这个在Spark里没有现成实现需要自己用mapPartitions写成本不低但对判断哪些特征是真正有用非常准。数据量大时也可以只在一个随机子样本上做。5. 从选完到上线VectorSlicer与Pipeline里那些容易翻车的细节特征选择算法本身跑通很容易真正麻烦的是把这个选择过程固化到训练和上线流程里保证训练、验证、线上推理三套环境用的是同一套特征逻辑。以下是我在工程化过程中反复踩过的坑。5.1 VectorAssembler VectorSlicer用列名代替索引很多教程里讲VectorSlicer都在用setIndices也就是按特征向量的位置切。但在真实项目里特征顺序会因为ETL改动、新特征插入而变化用死索引很容易选错特征。我强烈建议用setNames按特征名切片import org.apache.spark.ml.feature.VectorSlicer val slicer new VectorSlicer() .setInputCol(rawFeatures) .setOutputCol(selectedFeatures) .setNames(Array(age, income, click_rate)) val slicedDF slicer.transform(featureDF)用setNames的前提是你喂进来的rawFeatures是靠VectorAssembler的setInputCols组装出来的那些列名在元数据里有记录VectorSlicer才能正确解析。名字虽然比索引安全但也要注意如果两个不同阶段都往同一个向量里塞过特征列名可能被覆盖所以组装时尽量一口气完成不要多次嵌套。5.2 把特征选择过程固化进Pipeline并随模型一起保存我见过有同学把特征选择的结果存成一份特征名单然后用if else在代码里拼特征向量。这种方式短期能跑但一旦名单和训练代码不同步线上就出事。正确做法是把所有特征处理步骤包括StringIndexer、VectorAssembler、ChiSqSelector或树模型的筛选逻辑全部封装进Pipeline然后整体保存。pipelineModel.write.overwrite().save(/path/to/pipeline_model)上线推理时加载同一个pipelineModel直接transform线上数据Spark会保证训练时的索引映射、缺失值处理方式、选择逻辑完全一致。这样能避免绝大多数因为人工维护特征名单导致的不一致问题。5.3 三个我在生产环境踩过的坑未知类别、列缺失、driver端OOM第一个坑是未知类别。训练集里StringIndexer见过的类别列表有限线上数据一旦出现新类别默认行为是抛异常。不同Spark版本对handleInvalid的默认值还不一样2.x和3.x的默认行为有差异所以必须在每个StringIndexer上显式设置setHandleInvalid(keep)并在后续环节把这些unknown值统一归到一个桶里。第二个坑是列缺失。训练数据有的特征列线上数据表不一定有。VectorAssembler一旦发现输入列缺失直接抛IllegalArgumentException。我们在上线前会对齐特征Schema用Spark SQL把缺失列补成默认值保证训练和推理的列集合完全一致。第三个坑是driver端OOM。特征重要性、卡方统计量这些结果最终都要汇总到driver端。如果你在一个几亿行、几千列的数据集上直接把所有特征collect到driver画图内存基本瞬间打满。正确的做法是取样本画图比如df.sample(0.01)或者df.takeSample只把抽样后的结果回传。featureImportances本身已经是模型训练完汇总好的不需要额外collect全量数据。6. 特征集合的稳定性评估选出来的特征能扛多久特征选择做完、模型上线很多人就觉得事情结束了。但真实业务里特征的有效性是会随时间变化的。用户行为习惯会变、外部环境会变、数据口径会变上个月排名第一的特征下个月可能就失效了。所以我在团队里有一条硬性要求特征选择结果必须做稳定性评估并且定期监控。6.1 跨时间窗口算TopN集合的Jaccard重叠度最直观的稳定性评估办法把历史数据按时间切成多个窗口比如按周切每周单独跑一次ChiSqSelector或者树模型importance各取Top 20特征然后逐周计算集合的Jaccard重叠度。重叠度越高说明这套特征选择结果越稳定模型不容易因为数据分布波动而大幅变化。def jaccard(set_a, set_b): inter len(set_a set_b) union len(set_a | set_b) return inter / union if union 0 else 0正常情况下相邻两周的Top集合重叠度应该不低于0.6。如果连续几周重叠度都很低我会去排查两件事一是基础数据有没有发生口径变更二是某些特征本身是不是噪声太大需要做分箱、平滑或者干脆从候选名单里去掉。6.2 用PSI监控入选特征的分布漂移稳定性不仅看排名还要看分布。PSIPopulation Stability Index群体稳定性指数是风控和金融场景里常用的指标现在已经被广泛用在特征监控里。对每个入选特征把训练期的分布当作基准再用当前周期数据的分布做对比PSI Σ (实际占比 - 基准占比) × ln(实际占比 / 基准占比)PSI小于0.1说明特征分布基本稳定0.1到0.25之间说明有明显变化大于0.25说明该特征的分布已经严重漂移很可能不再适合继续作为模型的输入。PSI的计算在Spark里很容易实现先对特征做分箱再用group by统计每个箱子的占比都写在SQL里就好。6.3 把特征选择做成每周自动运行的特征体检报告最后我强烈建议把特征选择从一次性项目工作升级成常态化体检流程。具体做法是写一个Spark定时任务每周自动对最新一周的数据跑一次特征Report内容包括特征缺失率、常量列变化、卡方或importance排名、TopN集合重叠度、每个入选特征的PSI。结果写入一张Delta表再挂上预警规则一旦重叠度过低或者PSI超标就发告警。这个流程听起来工作量不小但实际用Pipeline跑一遍也就是一个Spark job的事情。当你的模型因为特征漂移出现问题时这份报告能帮你节省至少一整天排查时间也能让团队在讨论加特征、减特征时有一份客观依据而不是靠感觉投票。我在实际项目中还发现一个小技巧特征选择报告不要只记录最终选中的名单要把每次的完整排名也存下来。这样当某个低排名特征在某次报告中突然冲进前三时你能第一时间发现数据异动而不是等它真正带崩模型后才回头翻日志。特征选择这件事选完只意味着这个版本的工作结束真正的价值在于建立一套持续观察它的机制让每一次选择都经得起时间检验。
返回列表