【发布时间】:2018-06-14 19:01:01
【问题描述】:
我正在尝试在一些文本挖掘之后进行 kmean 聚类,但我找不到如何将 ParseWikipedia.termDocumentMatrix 的结果转换为 kmean.fit 方法所需的数据集中
scala> val (termDocMatrix, termIds, docIds, idfs) = ParseWikipedia.termDocumentMatrix(lemmas, stopWords, numTerms, sc)
scala> val kmeans = new KMeans().setK(5).setMaxIter(200).setSeed(1L)
scala> termDocMatrix.take(1)
res24: Array[org.apache.spark.mllib.linalg.Vector] = Array((1000,[32,166,200,223,577,645,685,873,926],[0.18132966949934762,0.3777537726516676,0.3178848913768969,0.43380819546465704,0.30604090845847254,0.46007361524957147,0.2076406414508386,0.2995665853335863,0.1742843713808876]))
scala> val modele = kmeans.fit(termDocMatrix)
<console>:66: error: type mismatch;
found : org.apache.spark.rdd.RDD[org.apache.spark.mllib.linalg.Vector]
required: org.apache.spark.sql.Dataset[_]
val modele = kmeans.fit(termDocMatrix)
我尝试了一些转换,但总是出错
scala> import spark.implicits._
import spark.implicits._
scala> val ss=org.apache.spark.sql.SparkSession.builder().getOrCreate()
scala> ss.createDataset(termDocMatrix)
<console>:67: error: Unable to find encoder for type stored in a Dataset. Primitive types (Int, String, etc) and Product types (case classes) are supported by importing spark.implicits._ Support for serializing other types will be added in future releases.
ss.createDataset(termDocMatrix)
和其他人(预期结果,因为它不是数据集)
val termDocRows = termDocMatrix.map(org.apache.spark.sql.Row(_))
val schemaVecteurs = StructType(Seq(StructField("features", VectorType, true)))
val termDocVectors = spark.createDataFrame(termDocRows, schemaVecteurs)
val termDocMatrixDense = termDocMatrix.map(e => e.toDense)
(并尝试 kmeans.fit 每个)。唯一给出不同错误的是 termDocVectors
val modele = kmeans.fit(termDocVectors)
18/01/05 01:14:52 ERROR Executor: Exception in task 0.0 in stage 560.0 (TID 1682)
java.lang.RuntimeException: Error while encoding: java.lang.RuntimeException: org.apache.spark.mllib.linalg.SparseVector is not a valid external type for schema of vector
if (assertnotnull(input[0, org.apache.spark.sql.Row, true]).isNullAt) null else newInstance(class org.apache.spark.ml.linalg.VectorUDT).serialize AS features#75
at org.apache.spark.sql.catalyst.encoders.ExpressionEncoder.toRow(ExpressionEncoder.scala:290)
有人有线索吗? 感谢您的帮助
另外测试后提供的线索:
我可以在哪里申请 DS ?
scala> termDocMatrix.toDS
<console>:69: error: value toDS is not a member of org.apache.spark.rdd.RDD[org.apache.spark.mllib.linalg.Vector]
termDocMatrix.toDS
使用元组...
我仍然有错误(这次不同)
val ds = spark.createDataset(termDocMatrix.map(Tuple1.apply)).withColumnRenamed("_1", "features")
ds: org.apache.spark.sql.DataFrame = [features: vector]
scala> val modele = kmeans.fit(ds)
java.lang.IllegalArgumentException: requirement failed: Column features must be of type org.apache.spark.ml.linalg.VectorUDT@3bfc3ba7 but was actually org.apache.spark.mllib.linalg.VectorUDT@f71b0bce.
最初的问题似乎已经解决了。现在我面临一个新的问题,因为我从 mllib.Rowmatrix 计算SVD,而 kmeans 似乎在等待 ml 向量。我只需要找到如何在 ml 包中计算 SVD...
【问题讨论】:
-
toDS怎么样 -
感谢您的建议...我完成了问题...
-
如果你不知道,这里有一些相关的documentation and sample code on SVD。
标签: scala apache-spark rdd apache-spark-dataset