【发布时间】:2015-08-25 14:28:55
【问题描述】:
我有一个 Apache spark 集群,其中包含一个主节点和三个工作节点。工作节点有 32 个核心和 124G 内存。我还在 HDFS 中有一个数据集,其中包含大约 6.5 亿条文本记录。这个数据集是一些像这样读入的序列化 RDD:
import org.apache.spark.mllib.linalg.{Vector, Vectors, SparseVector}
val vectors = sc.objectFile[(String, SparseVector)]("hdfs://mn:8020/data/*")
我想从这些记录中提取一百万个样本来做一些分析,所以我想试试val sample = vectors.takeSample(false, 10000, 0)。但是,这最终会失败并显示此错误消息:
15/08/25 09:48:27 ERROR Utils: Uncaught exception in thread task-result-getter-3
java.lang.OutOfMemoryError: Java heap space
at org.apache.spark.scheduler.DirectTaskResult$$anonfun$readExternal$1.apply$mcV$sp(TaskResult.scala:64)
at org.apache.spark.util.Utils$.tryOrIOException(Utils.scala:1239)
at org.apache.spark.scheduler.DirectTaskResult.readExternal(TaskResult.scala:61)
at java.io.ObjectInputStream.readExternalData(ObjectInputStream.java:1837)
at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:1796)
at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1350)
at java.io.ObjectInputStream.readObject(ObjectInputStream.java:370)
at org.apache.spark.serializer.JavaDeserializationStream.readObject(JavaSerializer.scala:69)
at org.apache.spark.serializer.JavaSerializerInstance.deserialize(JavaSerializer.scala:89)
at org.apache.spark.scheduler.TaskResultGetter$$anon$2$$anonfun$run$1.apply$mcV$sp(TaskResultGetter.scala:79)
at org.apache.spark.scheduler.TaskResultGetter$$anon$2$$anonfun$run$1.apply(TaskResultGetter.scala:51)
at org.apache.spark.scheduler.TaskResultGetter$$anon$2$$anonfun$run$1.apply(TaskResultGetter.scala:51)
at org.apache.spark.util.Utils$.logUncaughtExceptions(Utils.scala:1772)
at org.apache.spark.scheduler.TaskResultGetter$$anon$2.run(TaskResultGetter.scala:50)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1145)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:615)
at java.lang.Thread.run(Thread.java:745)
Exception in thread "task-result-getter-3" java.lang.OutOfMemoryError: Java heap space
at org.apache.spark.scheduler.DirectTaskResult$$anonfun$readExternal$1.apply$mcV$sp(TaskResult.scala:64)
at org.apache.spark.util.Utils$.tryOrIOException(Utils.scala:1239)
at org.apache.spark.scheduler.DirectTaskResult.readExternal(TaskResult.scala:61)
at java.io.ObjectInputStream.readExternalData(ObjectInputStream.java:1837)
at java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:1796)
at java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1350)
at java.io.ObjectInputStream.readObject(ObjectInputStream.java:370)
at org.apache.spark.serializer.JavaDeserializationStream.readObject(JavaSerializer.scala:69)
at org.apache.spark.serializer.JavaSerializerInstance.deserialize(JavaSerializer.scala:89)
at org.apache.spark.scheduler.TaskResultGetter$$anon$2$$anonfun$r
我知道我的堆空间用完了(我认为是在驱动程序上?),这是有道理的。执行hadoop fs -du -s /path/to/data,数据集在磁盘上占用了 2575 GB(但大小仅为约 850 GB)。
所以,我的问题是,我能做些什么来提取这个包含 1000000 条记录的样本(我稍后计划将其序列化到磁盘)?我知道我可以用较小的样本量做takeSample() 并在以后聚合它们,但我认为我只是没有设置正确的配置或做错了什么,这阻止了我以我想要的方式这样做。
【问题讨论】:
-
我最终不得不增加
spark.driver.memory和spark.driver.maxResultSize才能让事情正常进行。此外,根据接受的响应调整我的集群可能也有帮助。
标签: java scala apache-spark cloud