【发布时间】: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 的目标
每个数据样本在文件中每 2048 行出现一次(因此文件应该按顺序进行解析)
此代码需要在分布式集群上运行
我不需要耗尽内存 :)
提前致谢
【问题讨论】:
-
谢谢。如果是这样,@RohanAletty 提供的答案应该可以正常工作,只要您按索引对分组数据进行排序。
标签: scala apache-spark rdd