【问题标题】:Spark SQL in YARN: OutOfMemoryErrorYARN 中的 Spark SQL:OutOfMemoryError
【发布时间】:2021-10-13 09:26:16
【问题描述】:

问题描述

我正在尝试使用 Spark SQL 从存储在 HDFS 中的 ORC 文件中查询数据。当我处理小数据量(2-5 Gb)时,我不会遇到任何问题。但是,如果我尝试处理超过 400 Gb 的数据,我的执行程序会出错。

问题的根本原因是什么 - 驱动程序代码或集群配置?

21/08/09 14:28:35 ERROR executor.Executor: Exception in task 70.0 in stage 1.0 (TID 90) java.lang.OutOfMemoryError: Java heap space
        at org.apache.orc.impl.RecordReaderUtils.readDiskRanges(RecordReaderUtils.java:565)
        at org.apache.orc.impl.RecordReaderUtils$DefaultDataReader.readFileData(RecordReaderUtils.java:285)
        at org.apache.orc.impl.RecordReaderImpl.readAllDataStreams(RecordReaderImpl.java:1147)
        at org.apache.orc.impl.RecordReaderImpl.readStripe(RecordReaderImpl.java:1103)
        at org.apache.orc.impl.RecordReaderImpl.advanceStripe(RecordReaderImpl.java:1256)
        at org.apache.orc.impl.RecordReaderImpl.advanceToNextRow(RecordReaderImpl.java:1291)
        at org.apache.orc.impl.RecordReaderImpl.<init>(RecordReaderImpl.java:286)
        at org.apache.orc.impl.ReaderImpl.rows(ReaderImpl.java:669)
        at org.apache.spark.sql.execution.datasources.orc.OrcColumnarBatchReader.initialize(OrcColumnarBatchReader.java:130)
        at org.apache.spark.sql.execution.datasources.orc.OrcFileFormat.$anonfun$buildReaderWithPartitionValues$1(OrcFileFormat.scala:216)
        at org.apache.spark.sql.execution.datasources.orc.OrcFileFormat$$Lambda$1041/844883412.apply(Unknown Source)
        at org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1.org$apache$spark$sql$execution$datasources$FileScanRDD$$anon$$readCurrentFile(FileScanRDD.scala:116)
        at org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1.nextIterator(FileScanRDD.scala:169)
        at org.apache.spark.sql.execution.datasources.FileScanRDD$$anon$1.hasNext(FileScanRDD.scala:93)
        at org.apache.spark.sql.execution.FileSourceScanExec$$anon$1.hasNext(DataSourceScanExec.scala:503)
        at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.columnartorow_nextBatch_0$(Unknown Source)
        at org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.processNext(Unknown Source)
        at org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
        at org.apache.spark.sql.execution.WholeStageCodegenExec$$anon$1.hasNext(WholeStageCodegenExec.scala:755)
        at org.apache.spark.sql.execution.columnar.DefaultCachedBatchSerializer$$anon$1.hasNext(InMemoryRelation.scala:118)
        at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:458)
        at org.apache.spark.storage.memory.MemoryStore.putIterator(MemoryStore.scala:221)
        at org.apache.spark.storage.memory.MemoryStore.putIteratorAsValues(MemoryStore.scala:299)
        at org.apache.spark.storage.BlockManager.$anonfun$doPutIterator$1(BlockManager.scala:1423)
        at org.apache.spark.storage.BlockManager$$Lambda$585/222823930.apply(Unknown Source)
        at org.apache.spark.storage.BlockManager.org$apache$spark$storage$BlockManager$$doPut(BlockManager.scala:1350)
        at org.apache.spark.storage.BlockManager.doPutIterator(BlockManager.scala:1414)
        at org.apache.spark.storage.BlockManager.getOrElseUpdate(BlockManager.scala:1237)
        at org.apache.spark.rdd.RDD.getOrCompute(RDD.scala:384)
        at org.apache.spark.rdd.RDD.iterator(RDD.scala:335)
        at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
        at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:373)

我的 Spark 驱动程序代码和环境配置如下。

代码

我使用以下代码:

SparkSession spark = SparkSession.builder().getOrCreate();

spark.read()
     .option(MERGE_SCHEMA_OPTION, true)
     .orc("/path/to/folder/with/orc/group/by/partition")
     .persist(StorageLevel.MEMORY_AND_DISK()))
     .createOrReplaceTempView("viewName");

spark.sql("SELECT * FROM viewName WHERE KEY IN ('2021-08-01', '2021-08-02') AND time >= '2021-08-01T00:00:00.000Z' AND time <= '2021-08-02T23:59:59.000Z' AND code > 400")
      .write()
      .orc(properties.getResultLocation();

在我的解决方案中,我在 SQL 字符串中使用了所有运算符,例如 WHEREKEY IN。由于查询结构取决于许多参数,因此它是构建查询的最简单方法。

数据说明

数据显示为具有以下架构的 *.ORC 文件:

Field Type
message string
code int
rate double
time timestamp

所有带有数据的文件都存储在以时间值字段命名的分区文件夹中:

.
├── /path/to/folder/with/orc/group/by/partition
│   └──key=2021-08-01
│       ├── part_1.orc
│       ├── part_2.orc
│       └── part_...
│   └── key=2021-08-02
│       ├── part_1.orc
│       ├── part_2.orc
│       └── part_...
│   └── key=...
│       ├── ...
│       └── ...

出于测试目的,我使用了大约 400 GB 的数据(消耗了 1.2 TB 磁盘空间)。

环境配置

我在 YARN 集群中运行我的应用程序,其中包含 3 个节点管理器,配置如下:

RAM Cores
16 Gb 8

我提交带有属性的申请:

Property Value
spark.executors.cores 5
spark.executor.memory 14 Gb
spark.yarn.executor.memory.overhead 2 Gb
spark.driver.memory 14 Gb
spark.driver.cores 5
spark.executor.instances 2
spark.default.parallelism 20
spark.driver.extraJavaOptions -XX:+UseG1GC -XX:UseCompressedOops
spark.executor.extraJavaOptions -XX:+UseG1GC -XX:UseCompressedOops

Spark 版本 - 3.1.2 YARN 版本 - 3.0.0-cdh6.1.1

【问题讨论】:

    标签: java apache-spark apache-spark-sql hadoop-yarn


    【解决方案1】:

    通过压缩,您 400GB 的 orc 文件可以是 1.6TB 或更多的数据。此外,并行度为 20,因此每个执行程序可能会获得约 1/20 的数据(取决于可以下推多少查询以防止将数据读入内存) - 可能14GB的executor需要读取一个大于工作内存的chunk

    尝试增加执行器的并行度和/或内存(您可能还想添加执行器)。

    【讨论】:

      猜你喜欢
      • 2019-07-25
      • 2016-01-18
      • 2018-08-09
      • 2019-04-23
      • 2017-10-31
      • 2016-11-14
      • 1970-01-01
      • 2018-10-15
      • 1970-01-01
      相关资源
      最近更新 更多