【发布时间】: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 字符串中使用了所有运算符,例如 WHERE 和 KEY 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