You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.26 01:45:02