【发布时间】:2015-12-08 02:09:27
【问题描述】:
我尝试对来自两个大文本文件的两个数据集应用 join。两个文本文件都包含如下两列:
*col1* *document*
abc 1
aab 1
... ...
ccd 2
abc 2
... ...
我根据这两个文件的第一列加入这两个文件,并尝试找出文档有多少共同的 col1 值。两个文本文件的大小均为 10 GB。当我运行我的脚本时,spark 创建了 6 个阶段,每个阶段有 287 个分区。在这 6 个阶段中,有 4 个不同的阶段,一个 foreach 和一个 map。一切顺利,直到第五阶段的映射阶段。在那个阶段,火花停止处理分区,而是溢出到磁盘上,并且在溢出一万次后,它给出了与磁盘空间不足有关的错误。
我有 4 个内核和 8 GB 内存。我用 -Xmx8g 提供了所有内存。我也试过 set("spark.shuffle.spill", "true")。
My script:
{
...
val conf = new SparkConf().setAppName("ngram_app").setMaster("local[4]").set("spark.shuffle.spill", "false")
val sc = new SparkContext(conf)
val emp = sc.textFile("...doc1.txt").map { line => val parts = line.split("\t")
((parts(5)),parts(0))
}
val emp_new = sc.textFile("...doc2.txt").map { line => val parts = line.split("\t")
((parts(3)),parts(1))
}
val finalemp = emp_new.join(emp).
map { case((nk1) ,((parts1), (val1))) => (parts1 + "-" + val1, 1)}.reduceByKey((a, b) => a + b)
finalemp.foreach(println)
}
我应该怎么做才能避免大量溢出?
【问题讨论】:
标签: scala apache-spark diskspace