【问题标题】:Storagelevel in spark RDD MEMORY_AND_DISK_2() throw exceptionspark RDD MEMORY_AND_DISK_2() 中的存储级别抛出异常
【发布时间】:2016-06-20 09:30:58
【问题描述】:

谁能解释一下 rdd 的存储级别是如何工作的。

当我使用存储级别的持久方法时出现堆内存错误(StorageLevel.MEMORY_AND_DISK_2()) 但是,当我使用缓存方法时,我的代码可以正常工作。

根据 spark doc 文档使用默认存储级别 (MEMORY_ONLY) 缓存 Persist RDD。

我的代码出现堆错误

JavaRDD<String> rawData = sparkContext
                    .textFile(inputFile.getAbsolutePath())
                    .setName("Input File").persist(SparkToolConstant.rdd_stroage_level);

//          cache()

            String[] headers = new String[0];
            String headerStr = null;
            if (headerPresent) {
                headerStr = rawData.first();
                headers = headerStr.split(delim);
                List<String> headersList = new ArrayList<String>();
                headersList.add(headerStr);
                JavaRDD<String> headerRDD = sparkContext
                        .parallelize(headersList);
                JavaRDD<String> filteredRDD = rawData.subtract(headerRDD)
                        .setName("Raw data without header").persist(StorageLevel.MEMORY_AND_DISK_2());;
                rawData = filteredRDD;
            }

堆栈跟踪

 Job aborted due to stage failure: Task 0 in stage 3.0 failed 1 times, most recent failure: Lost task 0.0 in stage 3.0 (TID 10, localhost): java.lang.OutOfMemoryError: Java heap space
    at java.util.Arrays.copyOf(Arrays.java:2271)
    at java.io.ByteArrayOutputStream.grow(ByteArrayOutputStream.java:113)
    at java.io.ByteArrayOutputStream.ensureCapacity(ByteArrayOutputStream.java:93)
    at java.io.ByteArrayOutputStream.write(ByteArrayOutputStream.java:140)
    at java.io.BufferedOutputStream.flushBuffer(BufferedOutputStream.java:82)
    at java.io.BufferedOutputStream.write(BufferedOutputStream.java:126)
    at java.io.ObjectOutputStream$BlockDataOutputStream.drain(ObjectOutputStream.java:1876)
    at java.io.ObjectOutputStream$BlockDataOutputStream.setBlockDataMode(ObjectOutputStream.java:1785)
    at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1188)
    at java.io.ObjectOutputStream.writeObject(ObjectOutputStream.java:347)
    at org.apache.spark.serializer.JavaSerializationStream.writeObject(JavaSerializer.scala:44)
    at org.apache.spark.serializer.SerializationStream.writeAll(Serializer.scala:110)
    at org.apache.spark.storage.BlockManager.dataSerializeStream(BlockManager.scala:1176)
    at org.apache.spark.storage.BlockManager.dataSerialize(BlockManager.scala:1185)
    at org.apache.spark.storage.BlockManager.doPut(BlockManager.scala:846)
    at org.apache.spark.storage.BlockManager.putArray(BlockManager.scala:668)
    at org.apache.spark.CacheManager.putInBlockManager(CacheManager.scala:176)
    at org.apache.spark.CacheManager.getOrCompute(CacheManager.scala:79)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:242)
    at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:35)
    at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:277)
    at org.apache.spark.rdd.RDD.iterator(RDD.scala:244)
    at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:68)
    at org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:41)
    at org.apache.spark.scheduler.Task.run(Task.scala:64)
    at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:203)
    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)

Driver stacktrace:

Spark 版本:1.3.0

【问题讨论】:

  • 您能提供您的集群详细信息吗?
  • 我正在使用 spark-submit cmd 在本地模式、4 核和 2gb 驱动程序内存的本地系统上运行。
  • 工人内存?因为这个异常来自工人。
  • 如果我在本地模式(独立模式)下运行,我可以在哪里设置工作内存。
  • spark.executor.memory=2g

标签: java apache-spark


【解决方案1】:

看到这个问题很久没有得到答复,我发布此信息是为了提供一般信息,以及像我这样在此处搜索的人。

如果没有关于您的应用程序的更多细节,这类问题很难回答。一般来说,在序列化到磁盘时出现内存错误似乎是颠倒的。我建议你试试with Kryo serialization,如果你有很多额外的内存,请使用Alluxio (the software formerly known as Tachyon :) 进行“磁盘序列化”,这会加快速度。

更多来自 Spark 文档的 Tuning Data Storage, Serialized RDD Storage and (maybe helpful) GC Tuning

当您的对象仍然太大而无法有效存储时 这种调整,减少内存使用的更简单的方法是存储 它们以序列化的形式,使用序列化的 StorageLevels RDD persistence API,如 MEMORY_ONLY_SER。然后 Spark 将 将每个 RDD 分区存储为一个大字节数组。唯一的缺点 以序列化形式存储数据会降低访问时间,因为 动态反序列化每个对象。我们强烈推荐使用 Kryo 如果你想以序列化的形式缓存数据,因为它会导致很多 比 Java 序列化更小(当然也比原始 Java 对象)。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-09-05
    • 1970-01-01
    • 1970-01-01
    • 2015-10-11
    • 2019-02-24
    相关资源
    最近更新 更多