[
https://issues.apache.org/jira/browse/SPARK-20925?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Sean Owen resolved SPARK-20925.
-------------------------------
Resolution: Not A Problem
That doesn't mean the JVM is out of memory; it kind of means the opposite. It
thinks it can use more than YARN does, due to off-heap allocation. Setting the
heap size higher only helps if you make it so high that the default off-heap
cushion is sufficient. Increase spark.yarn.executor.memoryOverhead instead, as
your heap is likely far too big.
> Out of Memory Issues With org.apache.spark.sql.DataFrameWriter#partitionBy
> --------------------------------------------------------------------------
>
> Key: SPARK-20925
> URL: https://issues.apache.org/jira/browse/SPARK-20925
> Project: Spark
> Issue Type: Bug
> Components: SQL
> Affects Versions: 2.1.0
> Reporter: Jeffrey Quinn
>
> Observed under the following conditions:
> Spark Version: Spark 2.1.0
> Hadoop Version: Amazon 2.7.3 (emr-5.5.0)
> spark.submit.deployMode = client
> spark.master = yarn
> spark.driver.memory = 10g
> spark.shuffle.service.enabled = true
> spark.dynamicAllocation.enabled = true
> The job we are running is very simple: Our workflow reads data from a JSON
> format stored on S3, and write out partitioned parquet files to HDFS.
> As a one-liner, the whole workflow looks like this:
> ```
> sparkSession.sqlContext
> .read
> .schema(inputSchema)
> .json(expandedInputPath)
> .select(columnMap:_*)
> .write.partitionBy("partition_by_column")
> .parquet(outputPath)
> ```
> Unfortunately, for larger inputs, this job consistently fails with containers
> running out of memory. We observed containers of up to 20GB OOMing, which is
> surprising because the input data itself is only 15 GB compressed and maybe
> 100GB uncompressed.
> The error message we get indicates yarn is killing the containers. The
> executors are running out of memory and not the driver.
> ```Caused by: org.apache.spark.SparkException: Job aborted due to stage
> failure: Task 184 in stage 74.0 failed 4 times, most recent failure: Lost
> task 184.3 in stage 74.0 (TID 19110, ip-10-242-15-251.ec2.internal, executor
> 14): ExecutorLostFailure (executor 14 exited caused by one of the running
> tasks) Reason: Container killed by YARN for exceeding memory limits. 21.5 GB
> of 20.9 GB physical memory used. Consider boosting
> spark.yarn.executor.memoryOverhead.```
> We tried a full parameter sweep, including using dynamic allocation and
> setting executor memory as high as 20GB. The result was the same each time,
> with the job failing due to lost executors due to YARN killing containers.
> We were able to bisect that `partitionBy` is the problem by progressively
> removing/commenting out parts of our workflow. Finally when we get to the
> above state, if we remove `partitionBy` the job succeeds with no OOM.
--
This message was sent by Atlassian JIRA
(v6.3.15#6346)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]