【发布时间】:2015-01-05 11:05:07
【问题描述】:
我听很多人说 Spark 艺术擅长排序和分布式计算。目前,out 团队对 spark 和 scala 进行了一些研究。我们将在 Spark 上实现排序服务。现在,我已经设置了 spark 集群,并尝试在 spark 集群上运行和排序示例,但是排序的成本时间似乎很长。这是我的代码。
import org.apache.spark.{SparkConf, SparkContext}
import scala.collection.mutable.ListBuffer
import scala.util.Random
/**
* Created by on 1/1/15.
*/
object AdvancedSort {
/**
* bin/spark-submit --master spark://master:7077 --executor-memory 1024M --class com.my.sortedspark.AdvancedSort lib/sortedspark.jar 100000 3
* @param args
*/
def main(args: Array[String]) {
val sampleSize = if (args.length > 0) args(0).toInt else 100000
val slice = if (args.length > 1) args(1).toInt else 3
sort(sampleSize, slice)
}
def sort(listSize: Int, slice: Int): Unit = {
val conf = new SparkConf().setAppName(getClass.getName)
val spark = new SparkContext(conf)
val step1 = System.currentTimeMillis()
val data = genRandom(listSize)
val step2 = System.currentTimeMillis()
println(">>>>>>>>>> genRandom : " + (step2 - step1))
val distData = spark.parallelize(data, slice)
val step3 = System.currentTimeMillis()
println(">>>>>>>>>> parallelize : " + (step3 - step2))
val result = distData.sortBy(x => x, true).collect
val step4 = System.currentTimeMillis()
println(">>>>>>>>>> sortBy and collect: " + (step4 - step3))
println(">>>>>>>>>> total time : " + (step4 - step1))
printlnArray(result, 0, 10)
spark.stop()
}
/**
* generate random number
* @return
*/
def genRandom(listSize: Int): List[Int] = {
val range = 100000
var listBuffer = new ListBuffer[Int]
val random = new Random()
for (i <- 1 to listSize) listBuffer += random.nextInt(range)
listBuffer.toList
}
def printlnList(list: List[Int], start: Int, offset: Int) {
for (i <- start until start + offset) println(">>>>>>>>> list : " + i + " | " + list(i))
}
def printlnArray(list: Array[Int], start: Int, offset: Int) {
for (i <- start until start + offset) println(">>>>>>>>> list : " + i + " | " + list(i))
}
}
将以上代码部署到spark集群后,我在Master的Spark Home下运行如下命令:
bin/spark-submit --master spark://master:7077 --executor-memory 1024M --class com.my.sortedspark.AdvancedSort lib/sortedspark.jar 100000 3
以下是我最终得到的成本时间。
>>>>>>>>>> genRandom : 86
>>>>>>>>>> parallelize : 53
>>>>>>>>>> sortBy and collect: 6756
这看起来很奇怪,因为如果我在本地机器上通过 scala 的 sorted 方法运行 100000 个 Int 的随机数据,则成本时间会更快。
import scala.collection.mutable.ListBuffer
import scala.util.Random
/**
* Created by on 1/5/15.
*/
object ScalaSort {
def main(args: Array[String]) {
val list = genRandom(1000000)
val start = System.currentTimeMillis()
val result = list.sorted
val end = System.currentTimeMillis()
println(">>>>>>>>>>>>>>>>>> cost time : " + (end - start))
}
/**
* generate random number
* @return
*/
def genRandom(listSize: Int): List[Int] = {
val range = 100000
var listBuffer = new ListBuffer[Int]
val random = new Random()
for (i <- 1 to listSize) listBuffer += random.nextInt(range)
listBuffer.toList
}
}
scala 的 sorted 方法在本地机器上的花费时间
>>>>>>>>>>>>>>>>>> cost time : 169
在我看来,服装火花的排序时间有以下几个因素:
Master 和 Worker 之间的数据转换
在 Worker 上排序很快,通过合并可能很慢。
有没有火花大师知道为什么会这样?
【问题讨论】:
-
100000 个元素很小。正如您所说,由于开销,内存排序将击败并行版本。尝试一个适当大的数组。
-
嗨,Paul,尝试更大的数组将显示 spark 的优势。但我还有一个问题,是否可以将成本时间调整为 100 毫秒?
-
我不明白您所说的“调整成本时间”是什么意思。 Spark 是为解决大问题而设计的。如果你没有大问题,不要使用 Spark,或者当它没有更快时不要感到惊讶?
-
Chan 你的结果是一致的,考虑到你正在尝试做的事情,这并不奇怪。我同意@Paul - 如果您的数据大小适合单个节点的内存以进行排序(这里就是这种情况),那么为 Spark 设置基础设施的总成本是不值得的。只有当您可以摊销设置和分发成本时,您才会看到收益。
-
你需要测量大阵列,看看时间如何。开销将(相对)恒定,您应该会看到大型数组的加速。
标签: scala sorting apache-spark