代码之家  ›  专栏  ›  技术社区  ›  Igor Masternoy

KMeans ||用于Spark上的情绪分析

  •  2
  • Igor Masternoy  · 技术社区  · 10 年前

    我正在尝试编写基于Spark的情绪分析程序。为此,我使用了word2vec和KMeans聚类。从word2Vec,我在100维空间中收集了20k个单词/向量,现在我正在尝试对这个向量空间进行聚类。当我用默认的并行实现运行KMeans时,算法工作了3个小时!但在随机初始化策略下,这就像是8分钟。 我做错了什么?我有4个内核处理器和16 GB RAM的mac book pro机器。

    K~=4000 最大交互为20

    var vectors: Iterable[org.apache.spark.mllib.linalg.Vector] =
          model.getVectors.map(entry => new VectorWithLabel(entry._1, entry._2.map(_.toDouble)))
        val data = sc.parallelize(vectors.toIndexedSeq).persist(StorageLevel.MEMORY_ONLY_2)
        log.info("Clustering data size {}",data.count())
        log.info("==================Train process started==================");
        val clusterSize = modelSize/5
    
        val kmeans = new KMeans()
        kmeans.setInitializationMode(KMeans.K_MEANS_PARALLEL)
        kmeans.setK(clusterSize)
        kmeans.setRuns(1)
        kmeans.setMaxIterations(50)
        kmeans.setEpsilon(1e-4)
    
        time = System.currentTimeMillis()
        val clusterModel: KMeansModel = kmeans.run(data)
    

    spark上下文初始化如下:

    val conf = new SparkConf()
          .setAppName("SparkPreProcessor")
          .setMaster("local[4]")
          .set("spark.default.parallelism", "8")
          .set("spark.executor.memory", "1g")
        val sc = SparkContext.getOrCreate(conf)
    

    关于运行此程序的更新也很少。我在Intelij IDEA内部运行。我没有真正的Spark集群。但我想你的个人机器可以是Spark集群

    我看到程序挂在Spark代码LocalKMeans.scala的循环中:

    // Initialize centers by sampling using the k-means++ procedure.
        centers(0) = pickWeighted(rand, points, weights).toDense
        for (i <- 1 until k) {
          // Pick the next center with a probability proportional to cost under current centers
          val curCenters = centers.view.take(i)
          val sum = points.view.zip(weights).map { case (p, w) =>
            w * KMeans.pointCost(curCenters, p)
          }.sum
          val r = rand.nextDouble() * sum
          var cumulativeScore = 0.0
          var j = 0
          while (j < points.length && cumulativeScore < r) {
            cumulativeScore += weights(j) * KMeans.pointCost(curCenters, points(j))
            j += 1
          }
          if (j == 0) {
            logWarning("kMeansPlusPlus initialization ran out of distinct points for centers." +
              s" Using duplicate point for center k = $i.")
            centers(i) = points(0).toDense
          } else {
            centers(i) = points(j - 1).toDense
          }
        }
    
    2 回复  |  直到 10 年前
        1
  •  1
  •   CAFEBABE    10 年前

    使用初始化 KMeans.K_MEANS_PARALLEL 那就更复杂了 random 然而,它不应该有这么大的区别。我建议调查一下,是否是并行算法需要花费很多时间(它实际上应该比KMean本身更高效)。

    有关分析的信息,请参见: http://spark.apache.org/docs/latest/monitoring.html

    如果不是初始化占用了时间,则存在严重错误。然而,使用随机初始化不会对最终结果造成任何影响(只是效率更低!)。

    实际上,当你使用 KMeans.K_MEANS_PARALLEL表示K_MEANS平行 要初始化,您应该通过0次迭代获得合理的结果。如果不是这样的话,数据的分布可能会有一些规律,这会使KMean偏离轨道。因此,如果你没有随机分发数据,你也可以改变这一点。然而,这样的影响会让我惊讶地给出固定的迭代次数。

        2
  •  1
  •   Igor Masternoy    10 年前

    我在AWS上用3个从机(c3.xlarge)运行了spark,结果是一样的-问题是并行KMean在N次并行运行中初始化算法,但对于少量数据来说,它仍然非常慢,我的解决方案是继续使用随机初始化。 数据大小约为:4k个簇,21k个100维矢量。

    推荐文章