【发布时间】:2018-06-19 19:10:42
【问题描述】:
集群设置 -
Driver has 28gb
Workers have 56gb each (8 workers)
配置-
spark.memory.offHeap.enabled true
spark.driver.memory 20g
spark.memory.offHeap.size 16gb
spark.executor.memory 40g
我的工作 -
//myFunc just takes a string s and does some transformations on it, they are very small strings, but there's about 10million to process.
//Out of memory failure
data.map(s => myFunc(s)).saveAsTextFile(outFile)
//works fine
data.map(s => myFunc(s))
此外,我从我的程序中去集群/删除了 spark,它在具有 56gb 内存的单个服务器上完成得很好(成功保存到文件中)。这表明它只是一个火花配置问题。我查看了https://spark.apache.org/docs/latest/configuration.html#memory-management,我目前的配置似乎是我的工作需要更改的所有内容。我还应该改变什么?
更新-
数据 -
val fis: FileInputStream = new FileInputStream(new File(inputFile))
val bis: BufferedInputStream = new BufferedInputStream(fis);
val input: CompressorInputStream = new CompressorStreamFactory().createCompressorInputStream(bis);
br = new BufferedReader(new InputStreamReader(input))
val stringArray = br.lines().toArray()
val data = sc.parallelize(stringArray)
注意 - 这不会导致任何内存问题,即使它非常低效。我无法使用 spark 读取它,因为它会引发一些 EOF 错误。
myFunc,我不能真正发布它的代码,因为它很复杂。但基本上,输入字符串是一个分隔字符串,它会进行一些分隔符替换、日期/时间规范化等。输出字符串与输入字符串的大小大致相同。
此外,它适用于较小的数据大小,并且输出正确且与输入数据文件的大小大致相同。
【问题讨论】:
-
什么是
data,它是如何生成的?myFunc的定义是什么?myFunc中可能存在导致内存问题的问题,但更可能的问题是用于创建data的其他转换之一。不看代码就无法判断。所以我们需要从源文件中查看完整的流程。工作正常的那个这样做是因为它什么都不做。请记住,在您运行操作之前,Spark 中不会发生任何事情。 -
@puhlen 用信息更新了主帖
标签: scala apache-spark