Spark CountVectorizer返回UDT而非Vector的问题及转换需求
嘿,这个问题我在Spark 2.x版本折腾LDA的时候也碰到过,给你捋捋原因和解决办法:
为什么会显示成UDT?
首先别慌——Spark 2.3.0里CountVectorizer输出的features列本质就是org.apache.spark.ml.linalg.Vector类型,Databricks notebook里显示为UDT(User-Defined Type)只是因为Spark用UDT来封装MLlib的向量类,属于底层实现的显示方式,不是真的类型不对。
但如果你的LDA模型真的报错说类型不兼容,那大概率是你混用了Spark的两套机器学习API:
spark.ml:基于DataFrame的新API(官方推荐,后续会持续维护)spark.mllib:基于RDD的旧API(已经逐步废弃)CountVectorizer属于ml库,输出的是ml.linalg.Vector;如果你的LDA是用mllib里的实现,它需要的是mllib.linalg.Vector,这两个是完全不同的类,自然会报错。
怎么解决?
方案1:统一用spark.ml的API(最推荐)
Spark官方早就推荐用ml库的API了,ml里的LDA类直接支持CountVectorizer输出的向量类型,完全不需要额外转换。
举个完整的示例代码:
import org.apache.spark.ml.feature.CountVectorizer import org.apache.spark.ml.clustering.LDA // 假设你的数据集df有一列叫"tokens",是分词后的字符串数组 val countVectorizer = new CountVectorizer() .setInputCol("tokens") .setOutputCol("features") .setVocabSize(1000) // 根据你的需求调整词汇表大小 // 拟合向量器并转换数据 val cvModel = countVectorizer.fit(df) val featurizedData = cvModel.transform(df) // 用ml库的LDA训练模型 val lda = new LDA() .setK(5) // 主题数量 .setMaxIter(10) // 迭代次数 val ldaModel = lda.fit(featurizedData)
这样跑起来就不会有类型问题,因为从特征提取到模型训练全用的是ml生态的组件。
方案2:手动转换向量类型(迫不得已才用)
如果你因为历史代码或其他原因必须用mllib的LDA,那可以通过UDF把ml的向量转换成mllib的向量:
import org.apache.spark.mllib.clustering.LDA as MLlibLDA import org.apache.spark.mllib.linalg.Vector as MLlibVector import org.apache.spark.sql.functions.udf // 定义转换UDF,把ml向量转成mllib向量 val toMLlibVector = udf((vec: org.apache.spark.ml.linalg.Vector) => vec.toOld) // 生成mllib兼容的特征列 val mllibFeaturizedData = featurizedData .withColumn("mllib_features", toMLlibVector($"features")) .select($"mllib_features") .rdd.map(row => row.getAs[MLlibVector](0)) // 用mllib的LDA训练 val ldaModel = MLlibLDA.train(mllibFeaturizedData, k=5, maxIterations=10)
不过还是强烈建议你迁移到ml的API,因为mllib在后续Spark版本里会被彻底移除。
总结
你看到的UDT只是显示问题,核心问题是API混用。统一用spark.ml的组件就能直接用CountVectorizer的输出训练LDA,完全不需要额外转换类型。
内容的提问来源于stack exchange,提问作者Vince Robatel
相关产品推荐
相关产品推荐

