【问题标题】:converting RDD to vector with fixed length file data将 RDD 转换为具有固定长度文件数据的向量
【发布时间】:2016-03-06 04:13:14
【问题描述】:

我是 Spark + Scala 的新手,我的直觉仍在发展。我有一个包含许多数据样本的文件。每 2048 行代表一个新样本。我正在尝试将每个样本转换为向量,然后通过 k-means 聚类算法运行。数据文件如下所示:

123.34  800.18
456.123 23.16
...

当我使用非常小的数据子集时,我会从文件中创建一个 RDD,如下所示:

val fileData = sc.textFile("hdfs://path/to/file.txt")

然后使用以下代码创建向量:

val freqLineCount = 2048
val numSamples    = 200
val freqPowers = fileData.map( _.split(" ")(1).toDouble )
val allFreqs    = freqPowers.take(numSamples*freqLineCount).grouped(freqLineCount)
val lotsOfVecs  = allFreqs.map(spec => Vectors.dense(spec) ).toArray
val lotsOfVecsRDD = sc.parallelize( lotsOfVecs ).cache()

val numClusters = 2
val numIterations = 2
val clusters = KMeans.train(lotsOfVecsRDD, numClusters, numIterations)

这里的关键是我可以在一个字符串数组上调用.grouped,它会返回一个具有连续 2048 个值的数组数组。然后将其转换为向量并通过 KMeans 训练算法运行它是微不足道的。

我试图在更大的数据集上运行此代码并遇到java.lang.OutOfMemoryError: Java heap space 错误。大概是因为我在我的 freqPowers 变量上调用了take 方法,然后对该数据执行了一些操作。

牢记这一点,我将如何实现在此数据集上运行 KMeans 的目标

  1. 每个数据样本在文件中每 2048 行出现一次(因此文件应该按顺序进行解析)

  2. 此代码需要在分布式集群上运行

  3. 我不需要耗尽内存 :)

提前致谢

【问题讨论】:

  • 谢谢。如果是这样,@RohanAletty 提供的答案应该可以正常工作,只要您按索引对分组数据进行排序。

标签: scala apache-spark rdd


【解决方案1】:

你可以这样做:

val freqLineCount = 2048
val freqPowers = fileData.flatMap(_.split(" ")(1).toDouble)

// Replacement of your current code.
val groupedRDD = freqPowers.zipWithIndex().groupBy(_._2 / freqLineCount)
val vectorRDD = groupedRDD.map(grouped => Vectors.dense(grouped._2.map(_._1).toArray))

val numClusters = 2
val numIterations = 2
val clusters = KMeans.train(vectorRDD, numClusters, numIterations)

替换代码使用zipWithIndex() 和long 分割将RDD 元素分组为freqLineCount 块。分组后,将有问题的元素提取到它们的实际向量中。

【讨论】:

  • 您的缺失值排序。
  • 谢谢!根据the docs,“排序首先基于分区索引,然后是每个分区内的项目排序。”您是否碰巧知道分区索引是否保证与原始文件的顺序相同?
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2017-03-19
  • 2017-10-31
  • 1970-01-01
  • 2014-05-19
  • 1970-01-01
  • 2012-08-26
  • 2016-06-13
相关资源
最近更新 更多