[ 
https://issues.apache.org/jira/browse/SPARK-20925?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=16030108#comment-16030108
 ] 

Sean Owen commented on SPARK-20925:
-----------------------------------

This is better for the mailing list. Spark allocates off heap memory for lots 
of things (look up "spark tungsten"). Sometimes the default isn't enough. It's 
not a Spark issue per se, no, but a matter of how much YARN is asked to give 
the JVM. Partitioning isn't necessarily a trivial operation, and you might have 
some issue with key skew. By the way, the error message tells you about  
spark.yarn.executor.memoryOverhead already.

> 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]

Reply via email to