做算法开发做到一定阶段,你大概率会遇到这么个场景:喂给模型的特征从几十个一路涨到上千个,训练时间翻了几倍,线上效果却不升反降。我最近在Spark上重构一套用户复购预测的特征体系,核心工作之一就是特征选择。这篇文章把我在Apache Spark上做特征选择的这套流程完整梳理一遍,为什么选这套工具、参数怎么定、踩过哪些坑、线上怎么排查,都尽量说透。内容不止是API层面的罗列,更多是实操思路,适合正在用Spark调算法、搞特征工程的同行参考。
先说一个背景判断:Spark生态里特征选择相关的内置算法不像sklearn那么丰富,但胜在分布式,单机跑不动的数据量它能扛。很多人一上来就问"Spark有没有类似SelectKBest的东西",其实不光有,而且能和Pipeline、CrossValidator深度集成,关键是你得搞清楚每个工具的适用边界。下面我把这套东西完整拆开讲。
1. 特征选择在Spark算法开发中的位置
1.1 为什么单独把"特征选择"拎出来讲
在Spark上做算法开发,和单机实验有本质区别。单机上你用pandas加sklearn,几十万行数据几十个特征,怎么折腾都行。但一旦上了分布式,数据量千万级起步,特征维度上千,几件事会同时找上门:shuffle开销暴涨、模型训练时间不可控、特征之间共线性被放大、模型可解释性变差。这时候特征选择就不是"锦上添花"的预处理步骤,而是决定项目能不能按时交付的关键环节。
我在实际项目中见过不少团队,特征工程做完直接灌进树模型,靠feature importance事后筛一遍,效果也行,但在特定的业务场景下会吃亏。比如风控、精准营销这类需要向业务方解释"为什么用这几个特征"的项目,黑盒式筛选很难过业务评审。特征选择的价值就在于:它让你在建模前就压缩特征空间,同时保留业务可解释性,还能显著降低后续调参和训练成本。
另外还有一个非常现实的问题:Spark集群资源是按小时计费的。一个全量特征集跑一次十折交叉验证,和压缩到三分之一特征后跑,费用差距可能在一个数量级。特征选择省下的不只是时间,是真金白银的算力开销。
1.2 Filter、Wrapper、Embedded三种思路,Spark能覆盖到哪一层
特征选择的经典分类方法,无论单机还是分布式都适用,无非是Filter(过滤式)、Wrapper(包裹式)、Embedded(嵌入式)三种:
- Filter:先对每个特征独立打分,再按分数排序取TopN。典型代表是卡方检验、互信息、方差阈值。优点是不依赖模型、计算快;缺点是没考虑特征之间的交互。
- Wrapper:把特征子集本身当搜索空间,用模型效果来评价每个子集。典型代表是递归特征消除(RFE)。效果好,但计算代价极高,分布式场景下几乎不可行。
- Embedded:在模型训练过程中隐式完成特征选择。比如L1正则化、树模型的feature importance。这是目前工业界最常用的方式。
Spark MLlib原生支持的主要是Filter一层,也就是ChiSqSelector和FeatureHasher(FeatureHasher更准确说是特征变换/压缩,但常被用来做降维),配合VectorSlicer做规则化的手动筛选。Embedded层面,Spark的树模型(RandomForest、GBT)训练完可以直接取featureImportances,很多人拿它来二次筛特征,这是可行的,但它属于"训练之后"的选择,不是训练前的过滤。
我个人的项目经验是:不要迷信某一种方法。常见的高效套路是先用Filter快速压缩到几十个特征,再用树模型importance做第二轮确认。这样既能拿到分布式的高吞吐,又能兼顾特征交互信息。下文的所有实操都是围绕这条主线展开的。
2. Spark自带的特征选择工具箱
2.1 ChiSqSelector:主力中的主力
ChiSqSelector是Spark里最正统的特征选择器,基于卡方检验(Chi-Square Test)实现对每个特征与标签之间独立性的检验。它的核心思想是在问一个问题:这个特征在每个取值上的样本分布,和标签的正常分布是否显著不同?如果差异很大,说明这个特征和标签有相关性,值得保留;如果分布和期望几乎一样,那它对标签没有区分力,可以丢掉。
用生活化的类比来解释:你现在要判断"用户是否复购"(标签),有个特征是"是否领过优惠券"。如果领券用户的复购率和其他用户差不多,那这个特征对预测基本没用;如果领券用户复购率明显偏高,那它就有信息量。卡方检验就是量化"明显偏高"到底有多明显的工具。
在Spark中使用非常简单:
from pyspark.ml.feature import ChiSqSelector selector = ChiSqSelector( featuresCol="features", labelCol="label", outputCol="selectedFeatures", selectorType="fpr", fpr=0.05 ) model = selector.fit(df) result = model.transform(df)这里最关键的是selectorType参数,它决定了你按什么标准选特征。关于每种模式的具体策略,我在第三部分的实操里会给出参数建议,这里先记住一个结论:高维特征场景优先用fdr或fpr,不要一上来就固定numTopFeatures,具体原因后面讲。
2.2 VectorSlicer:规则优先时的手工切片
VectorSlicer是另一个非常实用的特征选择工具,它的定位和ChiSqSelector完全不同:不涉及任何统计检验,纯粹按位置或按名字切分特征向量。适合两种场景:一是你已经明确知道哪些特征要保留(比如业务上规定必须保留某些强特征),二是想快速验证某个特征子集的效果,不想重新组装特征。
按索引切:
from pyspark.ml.feature import VectorSlicer slicer = VectorSlicer( inputCol="features", outputCol="slicedFeatures", indices=[0, 2, 5] )按名称切:
slicer = VectorSlicer( inputCol="features", outputCol="slicedFeatures", names=["age", "order_cnt_30d", "city_rank"] )这里必须提醒一句:Spark 3.x的VectorSlicer里,indices和names不能同时指定,要么按位置、要么按名字。而且names模式要求inputCol对应的向量包含AttributeGroup信息,这个信息通常来自VectorAssembler设置或StringIndexer的输出。如果你直接手工创建了一个Vector,里面没有metadata,按名字切会直接报错。细节问题我在第四部分排查里详细说。
2.3 FeatureHasher:哈希变换与降维的组合拳
FeatureHasher严格来说不属于"选择",它做的是把类别特征映射到固定维度的稀疏向量,但它经常被用来解决一个和特征选择高度相关的问题——高基数类别特征导致维度爆炸。比如"用户ID"、"城市"、"商品类目"这类特征,如果做OneHot,一个特征就可能撑出上百万列,后续的卡方选择、模型训练都无从谈起。
FeatureHasher的思路是引入哈希函数,把每个类别值哈希后对numFeatures取模,映射到固定长度的向量里。用生活中的类比来说,OneHot是给每个类别单独开一扇门,特征多了门就多到离谱;哈希是只开固定数量的门,每个人凭哈希值往里"挤"。
from pyspark.ml.feature import FeatureHasher hasher = FeatureHasher( numFeatures=1024, inputCols=["city", "product_category", "user_segment"], outputCol="hashedFeatures" )哈希最经典的痛点是碰撞:不同类别值可能映射到同一个位置,造成信息混淆。实践中numFeatures不能设太小,否则碰撞率会快速上升。怎么估算合理值,我放到实操部分讲。
2.4 容易混淆的配角:VectorAssembler和PCA
有人在项目里会把PCA也当成特征选择用,这里要划清界限:PCA是特征提取(Feature Extraction),它把原始特征线性组合成新特征,虽然维度降了,但你保留的不再是原始特征,可解释性大打折扣。而真正的特征选择要求保留原始特征子集。在向业务方汇报"我们用了哪些特征"的时候,PCA出来的主成分是说不清楚的。
VectorAssembler则是所有特征处理步骤之前必须经过的关口。Spark里大部分特征变换器(包括ChiSqSelector、VectorSlicer)的输入要求是一个Vector类型的列,你需要用VectorAssembler把多个数值列拼成一个向量。
from pyspark.ml.feature import VectorAssembler assembler = VectorAssembler( inputCols=["age", "order_cnt_30d", "order_amt_30d", "category_cnt_30d"], outputCol="features" )这三个工具的配合关系是:VectorAssembler把分散字段拼成特征向量,ChiSqSelector在向量上做统计筛选,VectorSlicer做规则切片,FeatureHasher在高基数场景先行压缩。
3. 完整实操:从原始数据到最优特征子集
3.1 场景设定和数据准备
为了把流程讲清楚,我定义一个具体场景:电商用户复购预测。数据集大概1200万行,每行是一个用户,标签label表示一年内是否再次下单(二分类)。原始特征如下:
- gender:性别,字符串
- age:年龄,连续值
- city_rank:城市等级,连续值
- is_new_user:是否新客,0/1
- order_cnt_30d:近30天下单数,连续值
- order_amt_30d:近30天消费金额,连续值
- category_cnt_30d:近30天购买品类数,连续值
- coupon_used_30d:近30天用券次数,连续值
先加载数据并做基础清洗,把空值处理掉。特征选择的处理顺序非常讲究,缺失值必须在此之前解决,否则统计结果会被空值的分布干扰。
from pyspark.sql import SparkSession from pyspark.sql.functions import col, when spark = SparkSession.builder.appName("FeatureSelectionDemo").getOrCreate() df = spark.read.parquet("hdfs:///path/to/user_features.parquet") # 简易空值填充 df = df.fillna({ "age": 0, "order_cnt_30d": 0, "order_amt_30d": 0.0, "category_cnt_30d": 0, "coupon_used_30d": 0 })对于gender这种字符串特征,需要先转成数值。卡方检验本身处理的是离散分布,所以gender做StringIndexer加OneHotEncoder是常规操作,但要注意OneHotEncoder输出的向量列不能被VectorAssembler直接当作唯一输入,它本身就是一个稀疏向量列,需要单独展开或直接接上后续流程。这里为了简洁,我把gender先做StringIndexer转成索引数值,不作为单独的多列展开,避免喧宾夺主。
from pyspark.ml.feature import StringIndexer gender_indexer = StringIndexer(inputCol="gender", outputCol="gender_idx") df = gender_indexer.fit(df).transform(df)3.2 特征向量化与类型检查
接下来是VectorAssembler组装特征向量。这一步完成后,强烈建议先检查features列的Vector类型以及是否有非数值内容混进来。
numeric_cols = ["age", "city_rank", "is_new_user", "order_cnt_30d", "order_amt_30d", "category_cnt_30d", "coupon_used_30d", "gender_idx"] assembler = VectorAssembler(inputCols=numeric_cols, outputCol="features") df_assembled = assembler.transform(df)检查类型的标准姿势是看schema:
df_assembled.select("features").printSchema() # 期望输出: |-- features: vector (nullable = true)这里容易出问题的情况是:某些列虽然看起来是数值,但实际被读成了字符串,VectorAssembler会在这一步直接报"Data type StringType is not supported",需要提前用cast转换。
3.3 ChiSqSelector四种选择模式的参数策略
ChiSqSelector的selectorType有五个选项,我在项目里用下来,每一种都有它的适用语境:
| 选择模式 | 选择策略 | 适用场景 |
|---|---|---|
| numTopFeatures | 按卡方统计量降序取前N个 | 业务方明确要求"只保留20个特征"时 |
| percentile | 按卡方统计量降序取前百分之P | 只想压缩一个比例,不关心具体个数 |
| fpr | 控制假阳性率(默认0.05) | 希望保留统计显著的强相关特征 |
| fdr | Benjamini-Hochberg方法控制错误发现率 | 高维特征,允许多选一点再靠模型二次筛 |
| fwe | 控制家庭错误率(Bonferroni校正) | 严格要求少选错,宁可漏选也不误选 |
我的默认建议:不确定的时候用fdr=0.05。原因很简单,numTopFeatures和percentile完全不看统计显著性,可能把p值很高的弱相关特征也选进来,也可能把真正显著的少数特征漏掉。fpr在特征数量几千的时候会偏严,卡方检验在多重比较下容易放大假阳性;fwe又过于保守,往往筛完就没剩几个特征了。fdr是这两者之间的折中,在千万级数据上实际效果最好。
代码示例:
from pyspark.ml.feature import ChiSqSelector selector = ChiSqSelector( featuresCol="features", labelCol="label", outputCol="selectedFeatures", selectorType="fdr", fdr=0.05 ) selector_model = selector.fit(df_assembled)fit完成之后,可以通过selectedFeatures这个Attribute来查看最终到底保留了哪些原始特征。这一点很多教程没讲,但实际项目里非常重要,因为你总需要向业务方汇报"留下了哪些特征"。取法如下:
selected_attrs = selector_model.selectedFeatures # 这是一个列表,里面是Attribute对象,取名字可以用 .name feature_names = [attr.name for attr in selected_attrs if attr.name is not None] print(feature_names)注意这里的feature_names是否可获取,取决于之前VectorAssembler的inputCols是否通过setInputCols方式组装并携带元数据。如果读的是纯数据列没有attribute,这里拿到的就是null。
3.4 将特征选择嵌进Pipeline
很多人在特征选择上犯的第一个大错:在训练集上fit特征选择器,然后直接对测试集transform,但中间没有pipeline,导致流程在模型上线时要手工重放一遍。第二个大错更隐蔽:特征选择在大数据集上fit时把全量数据都用了,导致选择器提前看到了测试集分布,形成数据泄漏,后续交叉验证结果虚高。
正确的做法是把特征选择器和模型一起放进Pipeline,让Spark的交叉验证机制在每一折内部完成fit和transform,这样选择器只会看到训练折的数据,测试折的信息不会被"偷看"。
from pyspark.ml import Pipeline from pyspark.ml.classification import RandomForestClassifier rf = RandomForestClassifier(featuresCol="selectedFeatures", labelCol="label", maxDepth=8, numTrees=100) pipeline = Pipeline(stages=[assembler, selector, rf]) pipeline_model = pipeline.fit(train_df) test_pred = pipeline_model.transform(test_df)这里有一个容易被忽略的点:RandomForest的featuresCol要改成"selectedFeatures",而不是"features"。很多人写成默认的featuresCol="features",结果模型吃的还是全量特征,特征选择环节形同虚设。
另外,如果不做Pipeline、直接在原始df上调用selector.fit,随后再切分train/test,那么切分发生在特征选择之后,这个顺序是错的。正确顺序永远是:先切分,再在训练部分内部做选择。分布式环境下,建议用randomSplit切分后缓存训练集,减少重复计算。
3.5 用CrossValidator同时调特征选择和模型参数
特征选择的参数和模型参数在Spark里可以通过CrossValidator一起网格搜索。这算是把Spark MLlib的Pipeline能力用到极致了。以fpr参数为例,你可以在ParamGridBuilder里同时定义selector.fpr的候选值、rf.maxDepth的候选值,让交叉验证告诉我们哪组搭配最优。
from pyspark.ml.tuning import CrossValidator, ParamGridBuilder from pyspark.ml.evaluation import BinaryClassificationEvaluator grid = (ParamGridBuilder() .addGrid(selector.fpr, [0.01, 0.05, 0.1]) .addGrid(selector.selectorType, ["fpr", "fdr"]) .addGrid(rf.maxDepth, [6, 8, 10]) .addGrid(rf.numTrees, [50, 100]) .build()) evaluator = BinaryClassificationEvaluator(labelCol="label", metricName="areaUnderROC") cv = CrossValidator( estimator=pipeline, estimatorParamMaps=grid, evaluator=evaluator, numFolds=5, parallelism=4 ) cv_model = cv.fit(train_df)这种做法的计算代价不小,5折乘以几十组参数,在集群上跑起来以小时计。我的经验是先固定模型参数,单独搜特征选择参数,找到合理区间后再联合搜索,可以节省大量资源。另外,parallelism参数调高(比如4或8)能显著缩短总耗时,前提是集群有足够的executor。
搜索完成后,可以用cv_model.bestModel取出最优Pipeline里的features列和选择器信息,确认最终入选的特征集合。这一步的操作是层层取stage,比较绕,但确实能拿到结果。
best_pipeline = cv_model.bestModel best_selector = best_pipeline.stages[1] # 取决于Pipeline里的stage顺序 best_fpr = best_selector.getFpr() print("最优fpr:", best_fpr)4. 踩坑记录与排查技巧
4.1 ChiSqSelector要求非负特征值
这个报错几乎每个人都会遇到。ChiSqSelector在内部计算时是基于特征值构建列联表,特征值必须是零或正数。你只要有一个特征是负数,比如"消费金额差值"这类中心化处理过的特征,fit阶段就会直接抛异常:
requirement failed: ChiSqSelector requires non-negative feature values.解决思路有三个维度:
- 对负值特征做平移,比如MinMaxScaler归一化到[0,1]区间。
- 做分箱(Bucketizer),把连续值离散化到非负的桶ID。
- 对这个特征单独做绝对值或截断处理,视业务口径而定。
我倾向于优先用MinMaxScaler,因为它在保持分布形态的同时把范围压到0到1,卡方检验的结果也更稳定。注意MinMaxScaler也需要在Pipeline里,且要放在VectorAssembler之后、ChiSqSelector之前。
from pyspark.ml.feature import MinMaxScaler scaler = MinMaxScaler(inputCol="features", outputCol="scaledFeatures") selector = ChiSqSelector(featuresCol="scaledFeatures", ...)4.2 VectorSlicer按names切分时报错
按名称切特征向量,要求输入向量的AttributeGroup里包含特征名信息。如果你是用VectorAssembler从一堆DataFrame列拼出来的,特征是带metadata的,没问题。但如果你手工创建了向量,或者中间经过某些会丢失metadata的变换,就会看到类似"VectorSlicer requires attribute names"的异常。
排查方式很简单,看vector列的metadata:
df.select("features").schema[0].metadata如果metadata是空的,需要回到VectorAssembler重新组装,或者改为按indices切分。按索引切分不依赖metadata,但是你得维护一份"特征名到索引"的映射表,特征是可能随版本变化的,维护成本比names模式高。我的建议是能用names尽量用names,它可读性好、代码意图清晰。
4.3 连续特征直接丢给卡方检验会失真
卡方检验天然适用于离散特征。像"消费金额"这种连续特征,直接丢进ChiSqSelector,计算出来的统计量会把每个浮点值当作一个离散类别,导致严重过拟合——仿佛每个金额值都和标签独立相关,选择结果没有泛化意义。
实操中遇到连续特征,我会先做分箱。最简单的是用QuantileDiscretizer分成10到20个桶,让卡方检验在桶级别上计算。这样既保留了特征的整体形状,又让统计检验有实际意义。在Pipeline里的顺序是:
from pyspark.ml.feature import QuantileDiscretizer discretizer = QuantileDiscretizer( numBuckets=10, inputCol="order_amt_30d", outputCol="order_amt_30d_bucketed" )之后再进入VectorAssembler组装。QuantileDiscretizer的缺点是会在数据量大时产生较多shuffle,但对千万级数据来说可以接受,比单机方案靠谱得多。
4.4 特征哈希的碰撞问题
FeatureHasher看起来简单,埋的坑一点不少。假设你设置了numFeatures=1024,实际类别有5000个,那么平均每个桶要被5个类别共用,碰撞率会让模型学到的特征含义变得模糊。
估算碰撞影响的公式其实不复杂:如果总类别数远大于numFeatures,碰撞不可避免,但你真正要关心的是高频类别是否互相碰撞。比如"北京"和"上海"如果哈希到同一个位置,模型会把它们等同看待,这在城市特征上是致命的。解决方案有三个层级:
- 把numFeatures设大到至少是类别数的2到4倍。这里有一个经验判断:如果numFeatures低于类别总数的1/2,模型效果会明显下降。
- 对高基数特征分批处理,把最重要的前N个类别单独OneHot,剩余低频类别统一哈希到"其他"桶。
- 在哈希之前对特征值做标准化前缀,例如把城市值写成"city_北京",让不同类型特征不会因为哈希碰撞而混在一起,同时可以在最终交互中区分来源。
FeatureHasher输出的特征解释性较弱,如果业务方需要你解释具体某个特征的贡献,它可能不是一个好选择。我通常只在纯离线模型或基线模型里用它,最终上线模型尽量选可解释的特征集合。
4.5 数据泄漏:特征选择的顺序问题
这是我在团队评审里反复强调的一个问题。很多同学在Spark上写了一个独立的特征选择脚本,把全量数据跑一遍,选出一批特征,然后再切训练集、测试集,甚至基于这些特征做交叉验证。这样做出来的评估指标普遍虚高,因为选择器在全量数据上fit时已经"见过"测试集的标签分布了。
Spark的CrossValidator在Pipeline模式下能自动规避这个问题,因为每一折的fit都发生在该折的训练子集上。但如果你手动分折做实验,就必须自己保证:特征选择器的fit永远在训练折内部完成。换句话说,特征选择要和模型训练绑定在一个流水线里,而不是作为全局预处理步骤。
实战中一套合理的做法是:先用全量数据跑一次快速特征探索(比如看卡方值排序),确定一个大致的特征范围,然后在正式实验时把包含特征选择的Pipeline放进CrossValidator,用评估结果做最终判断。探索过程允许看全量数据,模型效果评估则必须严格走Pipeline。
5. 收尾:个人实操体会
最后分享一个我的习惯:特征选择做完之后,一定把入选特征的卡方p值、业务含义、模型重要性三张表对照着看一遍。卡方p值说明统计显著性,业务含义决定可解释性,树模型的feature importance则反映实际预测贡献。三者方向一致的特征,才是真正值得长期保留的特征。不一致的特征,大概率是数据质量或者样本分布有问题,需要回去查数据,而不是急着调参。
还有一个经验是:特征选择不是一劳永逸的。线上数据分布会漂移,上个月显著的强特征下个月可能就没用了。我现在的做法是每周跑一次卡方统计,按月更新特征稳定性报表,一旦发现某些特征的显著性和重要性出现大幅波动,就会触发特征集评审。这套机制比单次建模时的特征选择更有价值,它能帮你提前发现数据变化,而不是等模型效果掉了才回头排查。如果你的项目正卡在特征爆炸或者训练效率上,希望这篇能帮你少走几步弯路。