Spark Scala将Sparse稀疏特征转换为Dense稠密向量类型报错问询
报错原因
- 核心问题是Spark不同包下的Vector类型不匹配:你导入的是旧版
org.apache.spark.mllib.linalg.Vector(属于已经逐步废弃的mllib库类型),但spark.ml包下的CountVectorizer输出的feature列是新版org.apache.spark.ml.linalg.Vector类型,二者属于完全独立的类,即使结构体定义完全一致也无法自动转换,因此触发类型上抛失败的异常。
修复代码
仅需要替换导入的Vector包,同时修正你代码中列名拼接的语法问题即可正常运行:
import org.apache.spark.ml.Pipeline import org.apache.spark.ml.feature.{StringIndexer, OneHotEncoder} import org.apache.spark.ml.feature.CountVectorizer // 替换为ml包下的Vector,和CountVectorizer输出类型匹配 import org.apache.spark.ml.linalg.Vector import spark.implicits._ // 统计ocean_proximity的 distinct值数量 val distinctOceanProximities = dfRaw.select(col("ocean_proximity")).distinct().as[String].collect() val oceanProximityAsArrayDF = dfRaw.withColumn("ocean_proximity", array("ocean_proximity")) val countModel = new CountVectorizer().setInputCol("ocean_proximity").setOutputCol("feature").fit(oceanProximityAsArrayDF) val transformedDF = countModel.transform(oceanProximityAsArrayDF) transformedDF.show() def columnExtractor(idx: Int) = udf((v: Vector) => v(idx)) // 修正列名拼接的语法错误,把idx放到大括号内 val featureCols = (0 until distinctOceanProximities.size).map(idx => columnExtractor(idx)($"feature").as(s"${distinctOceanProximities(idx)}")) val toDense = udf((v:Vector) => v.toDense) val denseDF = transformedDF.withColumn("feature", toDense($"feature")) denseDF.show()
内容的提问来源于stack exchange,提问作者joesan
相关产品推荐
相关产品推荐

