Jeffrey Quinn created SPARK-20925:
-------------------------------------
Summary: 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.
--
This message was sent by Atlassian JIRA
(v6.3.15#6346)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]