PySpark如何将DenseVector数组转换为普通数组及实现加权投票分类器
问题解决办法
1. 报错与特征格式转换解决
你遇到的报错核心原因是接口不兼容:mlxtend的EnsembleVoteClassifier要求传入的基分类器符合Scikit-learn风格接口,必须实现predict()/predict_proba()方法,但你传入的是PySpark PipelineModel,它的预测接口是transform(),没有predict方法,所以触发报错。
首先解决DenseVector转标准Numpy数组的问题,直接调用PySpark DenseVector自带的toArray()方法即可:
# 训练集特征转换 y_pred3 = np.array([vec.toArray() for vec in pandas_df3['features']]) # 测试集特征转换 y_pred4 = np.array([vec.toArray() for vec in pandas_df4['features']])
接下来需要把PySpark PipelineModel包装成符合Sklearn接口的类,适配mlxtend的调用要求:
class PySparkClfWrapper: def __init__(self, spark_model, feature_col="features", pred_col="prediction"): self.model = spark_model self.feature_col = feature_col self.pred_col = pred_col def fit(self, X, y=None): # 你的s1、s2已经预训练完成,且refit设为False,此处直接返回自身即可 return self def predict(self, X): # 把numpy输入转成PySpark可用的DataFrame from pyspark.ml.linalg import Vectors spark_input = spark.createDataFrame( [(Vectors.dense(x),) for x in X], schema=[self.feature_col] ) # 调用PySpark模型预测并提取结果 pred_result = self.model.transform(spark_input) return np.array([row[self.pred_col] for row in pred_result.collect()])
初始化集成模型时用包装类包裹两个PySpark模型即可:
eclf = EnsembleVoteClassifier( clfs=[PySparkClfWrapper(s1), PySparkClfWrapper(s2)], weights=[1,1], refit=False )
如果需要用软投票,只需要在包装类中新增predict_proba方法,提取模型输出的probability列返回即可。
2. PySpark原生加权投票分类器实现
PySpark没有提供和EnsembleVoteClassifier完全同名的现成接口,但可以基于PySpark原生API实现加权投票,不需要转成Pandas处理,更适合大数据场景:
硬投票实现示例
from pyspark.sql import functions as F from pyspark.sql.types import DoubleType # 定义两个模型的权重 weight_1 = 1 weight_2 = 1 # 先给数据集加唯一标识,方便后续合并预测结果 qa1 = qa1.withColumn("sample_id", F.monotonically_increasing_id()) # 分别用两个模型预测 pred1 = s1.transform(qa1).select("sample_id", F.col("prediction").alias("pred1")) pred2 = s2.transform(qa1).select("sample_id", F.col("prediction").alias("pred2")) # 合并两个模型的预测结果 merge_pred = pred1.join(pred2, on="sample_id", how="inner") # 自定义加权投票UDF @F.udf(returnType=DoubleType()) def vote_cal(p1, p2): vote_map = {} vote_map[p1] = vote_map.get(p1, 0) + weight_1 vote_map[p2] = vote_map.get(p2, 0) + weight_2 return max(vote_map.items(), key=lambda x:x[1])[0] # 生成最终预测结果 final_result = merge_pred.withColumn("final_pred", vote_cal(F.col("pred1"), F.col("pred2")))
如果需要软投票,只需要把上述逻辑中的prediction列替换为probability列,按权重对概率加权求和后取概率最高的类别即可。
内容的提问来源于stack exchange,提问作者Anu
相关产品推荐
相关产品推荐

