【问题标题】:Spark: java.lang.IllegalArgumentException: requirement failed kmeans (mllib)Spark:java.lang.IllegalArgumentException:要求失败的kmeans(mllib)
【发布时间】:2018-10-30 10:55:48
【问题描述】:

我正在尝试使用 kmeans 进行聚类应用。 我的数据集是: https://archive.ics.uci.edu/ml/datasets/ElectricityLoadDiagrams20112014#

我对 spark 没有太多经验,我只工作了几个月,当我尝试应用具有输入的 kmean.train 时发生错误:向量、num_cluster 和迭代。

我在本地运行,是不是我的机器不能计算这么多数据?

主要代码为:

import org.apache.spark.sql.SparkSession
import org.apache.spark.SparkConf
import org.apache.spark.SparkContext
import scala.collection._
import org.apache.spark.sql.functions._
import org.apache.spark.sql.functions.udf
import org.apache.spark.sql.Row

import org.apache.spark.mllib.linalg.{Vector, Vectors}
import org.apache.spark.mllib.clustering.{KMeans, KMeansModel}

object Preprocesado {
def main(args: Array[String]) {
val spark = SparkSession.builder.appName("Preprocesado").getOrCreate()
import spark.implicits._ 
    val sc = spark.sparkContext 

    val datos = spark.read.format("csv").option("sep", ";").option("inferSchema", "true").option("header", "true").load("input.csv")
var df= datos.select("data", "MT_001").withColumn("data", to_date($"data").cast("string")).withColumn("data", concat(lit("MT_001 "), $"data"))
val col=datos.columns

for(a<- 2 to col.size-1) {
      var user = col(a)
      println(user)
      var df_$a = datos.select("data", col(a)).withColumn("data", to_date($"data").cast("string")).withColumn("data", concat(lit(user), lit(" "),  $"data"))
       df = df.unionAll(df_$a)
      }

val rd=df.withColumnRenamed("MT_001", "values")
val df2 = rd.groupBy("data").agg(collect_list("values"))

val convertUDF = udf((array : Seq[Double]) => {
  Vectors.dense(array.toArray)
})
val withVector = df2.withColumn("collect_list(values)", convertUDF($"collect_list(values)"))
val items : Array[Double] = new Array[Double](96)
val vecToRemove = Vectors.dense(items)
def vectors_unequal(vec1: Vector) = udf((vec2: Vector) => !vec1.equals(vec2))
val filtered = withVector.filter(vectors_unequal(vecToRemove)($"collect_list(values)"))
val Array(a, b) = filtered.randomSplit(Array(0.7,0.3))
val trainingData = a.select("collect_list(values)").rdd.map{x:Row => x.getAs[Vector](0)}
val testData = b.select("collect_list(values)").rdd.map{x:Row => x.getAs[Vector](0)}
trainingData.cache()
testData.cache()
val numClusters = 4
val numIterations = 20
val clusters = KMeans.train(trainingData, numClusters, numIterations)
clusters.predict(testData).coalesce(1,true).saveAsTextFile("output")


spark.stop()
  }
}

当我编译时没有错误。 然后我提交:

spark-submit \
  --class "spark.Preprocesado.Preprocesado" \
  --master local[4] \
  --executor-memory 7g \
  --driver-memory 6g \
  target/scala-2.11/preprocesado_2.11-1.0.jar

问题出在集群中: 这是错误:

18/05/20 16:45:48 ERROR Executor: Exception in task 10.0 in stage 7.0 (TID 6347)
java.lang.IllegalArgumentException: requirement failed
    at scala.Predef$.require(Predef.scala:212)
    at org.apache.spark.mllib.util.MLUtils$.fastSquaredDistance(MLUtils.scala:486)
    at org.apache.spark.mllib.clustering.KMeans$.fastSquaredDistance(KMeans.scala:589)
    at org.apache.spark.mllib.clustering.KMeans$$anonfun$findClosest$1.apply(KMeans.scala:563)
    at org.apache.spark.mllib.clustering.KMeans$$anonfun$findClosest$1.apply(KMeans.scala:557)
    at scala.collection.immutable.List.foreach(List.scala:381)
    at org.apache.spark.mllib.clustering.KMeans$.findClosest(KMeans.scala:557)
    at org.apache.spark.mllib.clustering.KMeans$.pointCost(KMeans.scala:580)
    at org.apache.spark.mllib.clustering.KMeans$$anonfun$initKMeansParallel$2.apply(KMeans.scala:371)
    at org.apache.spark.mllib.clustering.KMeans$$anonfun$initKMeansParallel$2.apply(KMeans.scala:370)
    at scala.collection.Iterator$$anon$11.next(Iterator.scala:409)
    at org.apache.spark.storage.memory.MemoryStore.putIteratorAsValues(MemoryStore.scala:216)
    at org.apache.spark.storage.BlockManager$$anonfun$doPutIterator$1.apply(BlockManager.scala:1038)
    at org.apache.spark.storage.BlockManager$$anonfun$doPutIterator$1.apply(BlockManager.scala:1029)
    at org.apache.spark.storage.BlockManager.doPut(BlockManager.scala:969)
    at org.apache.spark.storage.BlockManager.doPutIterator(BlockManager.scala:1029)
    at org.apache.spark.storage.BlockManager.getOrElseUpdate(BlockManager.scala:760)
    at org.apache.spark.rdd.RDD.getOrCompute(RDD.scala:334)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:285)
    at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:87)
    at org.apache.spark.scheduler.Task.run(Task.scala:108)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:338)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)

我该如何解决这个错误?

谢谢

【问题讨论】:

  • 这个特定的错误是由输入数据的维度不一致(输入向量的长度不同)引起的。与错误无关,但使用 rd.groupBy("data").agg(collect_list("values")) 创建功能看起来不正确。
  • 感谢您的回答,我检查了输入向量(trainingData)的长度,并且都具有相同的长度。 trainingData.foreach(x =&gt; println(x.size))
  • 如果您不需要使用集群,而是在本地运行,那么请考虑使用 Spark 的替代方案。它们可能会快很多,尤其是对于kmeans。

标签: scala apache-spark cluster-analysis


【解决方案1】:

我认为您正在以错误的方式生成 DataFrame df 并因此生成 df2。 也许您正在尝试这样做:

case class Data(values: Double, data: String)
var df = spark.emptyDataset[Data]

df = datos.columns.filter(_.startsWith("MT")).foldLeft(df)((df, c) => {
                val values = col(c).cast("double").as("values")
                val data = concat(lit(c), lit(" "), to_date($"_c0").cast("string")).as("data")

                df.union(datos.select(values, data).as[Data]) 
    })


val df2 = df.groupBy("data").agg(collect_list("values"))

我认为,您只需要两列:datavalues,但在 for 循环中,您将生成一个包含 140256 列的 DataFrame(每个属性一个列)也许这就是你问题的根源。

pd:对不起我的英语!

【讨论】:

    猜你喜欢
    • 2015-08-24
    • 2017-08-24
    • 2019-05-28
    • 2018-08-30
    • 1970-01-01
    • 2023-03-28
    • 1970-01-01
    • 2014-12-11
    • 2018-11-11
    相关资源
    最近更新 更多