Most Spark memory problems are diagnosed backwards. A job fails, someone doubles spark.executor.memory, the job passes, and the cluster now runs half as many executors as it could. The next failure gets the same treatment. Tuning goes better when you start from the symptom, identify which of a handful of memory regions actually ran out, and change the one setting that controls that region.

This article is that symptom-first workflow. It assumes you know the unified memory model in outline; the memory management article explains execution and storage borrowing in detail. Here the focus is practical: the anatomy of an executor container, how to tell heap OOM, container kills, driver OOM and heavy spill apart, a worked sizing example for a real node, per-task memory and partition sizing, PySpark and off-heap settings, and a checklist.

Anatomy of an executor container

An executor is a JVM running inside a container that a cluster manager limits. The container limit is the sum of four settings. The Spark configuration reference states that the maximum container memory is the sum of spark.executor.memoryOverhead, spark.executor.memory, spark.memory.offHeap.size and spark.executor.pyspark.memory.

Inside the heap, Spark reserves 300 MB, then gives spark.memory.fraction (default 0.6) of the rest to a unified pool shared by execution memory, used by shuffles, joins, sorts and aggregations, and storage memory, used by cached data and broadcast variables. spark.memory.storageFraction (default 0.5) is the part of the pool where cached blocks are immune to eviction. The remaining 40 percent is user memory: your objects, UDF state and Spark's internal metadata, which Spark does not track.

Overhead covers everything the JVM and its neighbours use outside the heap: metaspace, thread stacks, direct buffers used by the network layer, native libraries and, unless spark.executor.pyspark.memory is set, Python worker processes. Its default is spark.executor.memoryOverheadFactor (0.10, since Spark 3.3) times executor memory, with a 384 MiB minimum. On Kubernetes the factor defaults to 0.40 for non-JVM jobs, because Python and R workers need much more native memory.

One executor container, from the outside in (worked example: 33g heap, 5 cores)Container limit enforced by YARN or the kubelet: about 36.3 GBJVM heap: spark.executor.memory = 33gReserved 300 MBStorageprotected half, about 9.8 GBExecutionshuffle, join, sort, aggUnified pool = (heap - 300 MB) x 0.6, about 19.6 GB; the boundary movesUser memory, about 13.1 GBmemoryOverhead3.3 GB: metaspace, threads,netty, native libs, PythonoffHeap.size0 here; added if setpyspark.memoryunset here; added if setHeap OOM = a box inside the blue area is too small. Container kill = the sum of everything exceeds the outer limit.
The executor container for the worked example below. The four outer settings add up to the limit the cluster manager enforces.

Four symptoms, four regions

Each failure points to one region. Identify it before changing anything.

SymptomRegion that ran outFirst lever
Executor log shows OutOfMemoryError: Java heap spaceHeap: execution, user memory or one huge recordMore partitions, fix skew, then executor memory
YARN kills the container for exceeding physical memory; Kubernetes reports OOMKilled, exit code 137Everything outside the heap: overhead, Python, nativememoryOverhead or pyspark.memory
Driver OOM, or a job aborted for exceeding spark.driver.maxResultSizeDriver heap, from collect or large broadcastsStop collecting; raise driver memory last
No failure, but high GC time and large spillExecution memory per taskPartition count, cores per executor

The second row is the most often misread. When the cluster manager kills a container, the JVM never throws, so there is no heap dump and the executor log simply ends. The YARN message mentions exceeding physical memory limits and suggests boosting spark.executor.memoryOverhead; on Kubernetes the pod status shows OOMKilled. Raising spark.executor.memory in response barely helps: 4 GB more heap buys about 0.4 GB more overhead with the default factor, or nothing if overhead is set explicitly, and fewer executors now fit per node.

The fourth row is not a failure, but it costs more. The stage page of the Spark UI shows Spill (Memory) and Spill (Disk) per task, and the Executors tab shows GC time. Spill means execution memory per task was too small for the data each task handled; spilling is correct behaviour, but heavy spill multiplies disk I/O and runtime.

Worked example: sizing executors for a 128 GB node

Take worker nodes with 16 cores and 128 GB of RAM, where YARN is allowed 112 GB per node and one core is left for the operating system and daemons. A common starting shape is about five cores per executor, which keeps HDFS and shuffle client concurrency reasonable and gives three executors per node.

Divide memory the same way: 112 GB over three executors is about 37.3 GB per container. Overhead at 10 percent of heap means the heap is roughly 37.3 divided by 1.1, so set spark.executor.memory=33g, which gives 3.3 GB of overhead and a 36.3 GB container. YARN rounds container requests up to a multiple of its minimum allocation, so leave a little slack rather than filling the node exactly.

Now compute what each task gets. The heap is 33,792 MB. Unified memory is (33,792 minus 300) times 0.6, about 20,095 MB (19.6 GiB), so the protected storage half is about 9.8 GiB. User memory is 13,397 MB, about 13.1 GiB. Execution memory is shared among running tasks: with N active tasks, each task is guaranteed at least 1/(2N) of the execution pool and can use up to 1/N of it before it must spill. With five tasks and nothing cached, each task can count on about 2 GB and can reach about 3.9 GB. That number, not the 33 GB heap, is what a single sort or hash aggregation actually works with.

spark-submit \
  --conf spark.executor.instances=30 \
  --conf spark.executor.cores=5 \
  --conf spark.executor.memory=33g \
  --conf spark.executor.memoryOverhead=3400m \
  --conf spark.driver.memory=8g \
  --conf spark.sql.adaptive.enabled=true \
  --conf spark.sql.adaptive.advisoryPartitionSizeInBytes=128m \
  --conf spark.sql.shuffle.partitions=2000 \
  job.py

Setting overhead explicitly, as above, makes the container size visible in the job definition rather than derived from a factor someone may change later.

Making tasks fit: partitions, skew and cores

Since per-task execution memory is fixed by the executor shape, the main way to make a task fit is to give it less data. For shuffles, that means more shuffle partitions. The default spark.sql.shuffle.partitions=200 is small for large jobs: 200 partitions of a 1 TB shuffle are 5 GB each, far beyond the 2 to 4 GB each task can use. Set a generous initial number and let adaptive query execution coalesce small partitions after the shuffle, towards spark.sql.adaptive.advisoryPartitionSizeInBytes.

Skew defeats this. If one key holds 40 GB, one task receives 40 GB whatever the partition count. The heap OOM then comes from one straggler task while its siblings finish quickly. Look at the task duration and shuffle read distribution on the stage page: a maximum far above the median is skew. Enable AQE skew-join handling, or salt the key; the data skew article covers both.

Cores per executor are the other lever. Fewer concurrent tasks per executor means a larger share of execution memory for each. Dropping from five cores to four raises each task's ceiling from one fifth to one quarter of the pool, at the cost of parallelism. Use this for a memory-hungry stage when more partitions are not possible, for example a window function over a large group.

Two more causes of heap OOM are not about execution memory at all. A single enormous record, such as a large array column or a JSON document, must fit in memory as one object, and no partitioning helps. And user memory overflows when a UDF or mapPartitions function accumulates state, for example by building a list of the whole partition. Stream through iterators instead.

PySpark and off-heap memory

PySpark adds processes outside the JVM. Each Python worker that runs UDFs, pandas UDFs or Arrow conversion lives in the container and counts against its limit, but not against the heap. A pandas UDF materialises a whole Arrow batch as a pandas DataFrame, and pandas operations often make copies, so peak Python memory can be several times the batch size. spark.sql.execution.arrow.maxRecordsPerBatch (default 10,000 rows) caps the batch; lower it for wide rows.

If Python memory is not budgeted, it silently comes out of overhead and causes container kills. Setting spark.executor.pyspark.memory gives Python its own budget, which is added to the container request on YARN and Kubernetes and enforced as a limit on the workers. On Kubernetes the higher 0.40 default overhead factor for non-JVM jobs exists for the same reason.

Off-heap memory, enabled with spark.memory.offHeap.enabled and sized with spark.memory.offHeap.size, lets Tungsten keep execution data outside the heap and reduces GC pressure for large aggregations. The documentation is explicit that this setting has no effect on heap size, so shrink spark.executor.memory by the same amount if the container must stay the same size.

Measuring before and after

Tune from measurements. The monitoring REST API exposes per-executor peak memory metrics, and per-stage spill totals. The script below lists the worst stages by spill and each executor's peak heap and execution memory, which together tell you whether you are close to a limit or comfortably below it.

import requests

UI = "http://driver-host:4040/api/v1"     # or the history server
app = requests.get(f"{UI}/applications").json()[0]["id"]

stages = requests.get(f"{UI}/applications/{app}/stages").json()
for s in sorted(stages, key=lambda s: s["diskBytesSpilled"], reverse=True)[:5]:
    print(f"stage {s['stageId']}: spill disk {s['diskBytesSpilled'] / 2**30:.1f} GiB, "
          f"memory {s['memoryBytesSpilled'] / 2**30:.1f} GiB, tasks {s['numTasks']}")

for e in requests.get(f"{UI}/applications/{app}/executors").json():
    peak = e.get("peakMemoryMetrics") or {}
    gc_share = e["totalGCTime"] / max(e["totalDuration"], 1)
    print(f"exec {e['id']}: heap peak {peak.get('JVMHeapMemory', 0) / 2**30:.1f} GiB, "
          f"exec mem peak {peak.get('OnHeapExecutionMemory', 0) / 2**30:.1f} GiB, GC {gc_share:.0%}")

Container-level resident memory, including Python workers, is reported only when process-tree metrics are enabled with spark.executor.processTreeMetrics.enabled. Turn it on while investigating container kills; it shows whether the JVM or Python is growing.

A GC share above about 10 percent of task time is a sign that the heap is crowded. Before changing collectors, check whether cached data you no longer need is sitting in storage memory, and unpersist it.

Driver memory

The driver is a single JVM that plans queries, tracks tasks and receives results. It runs out of memory for three reasons: collect() or toPandas() pulling a large result back; building a large broadcast variable or broadcast join table, which is first assembled on the driver; and very large plans with tens of thousands of tasks or files, whose metadata adds up. spark.driver.maxResultSize (default 1g) aborts a job whose serialised results exceed the limit, which is a guard rail, not a problem to remove. Write large results to storage instead, and lower spark.sql.autoBroadcastJoinThreshold if broadcasts are the cause.

Failure modes in tuning itself

Common tuning mistakes, and what they cost:

  • Raising heap for a container kill. The container grows by heap plus 10 percent while native use stays the same, and fewer executors fit per node. Raise overhead or budget Python instead.
  • One giant executor per node. A 100 GB heap has long GC pauses, and many tasks contend for one pool. Several mid-sized executors are usually better.
  • Raising spark.memory.fraction to stop spills. It takes from user memory, so a job that was spilling now fails with heap OOM in a UDF. The documentation recommends leaving it at the default.
  • Caching everything. Cached data occupies storage memory and squeezes execution; cache only data that is reused, and unpersist it after.
  • Tuning for the average task. Memory fails on the largest task. Read maxima, not medians.

Trade-offs

Bigger executors share broadcast tables and cached data across more tasks and reduce per-executor overhead, but they suffer longer GC pauses and lose more work when one dies. Smaller executors isolate failures and pause briefly, but duplicate broadcasts in every JVM. More shuffle partitions reduce per-task memory but add scheduling overhead and small files. Off-heap memory cuts GC time but must be budgeted explicitly. Spilling is slower but always safer than failing, so a little spill on the largest stage is an acceptable result.

What to do next

  1. Classify your last memory failure using the symptom table before changing any setting.
  2. Write down the container formula for your jobs: heap plus overhead plus off-heap plus PySpark memory, against the node and cluster-manager limits.
  3. Compute per-task execution memory for your executor shape, as in the worked example.
  4. Size shuffle partitions so that the largest partition fits in that per-task budget, and enable AQE.
  5. Check the stage page for skew and spill on the three slowest stages.
  6. For PySpark jobs, set spark.executor.pyspark.memory and lower the Arrow batch size for wide rows.
  7. Collect peak memory metrics from the REST API for a week and tune from the maxima.
  8. Read the Spark on Kubernetes article if you are moving off YARN, since the overhead defaults differ there.
Key takeaway: Spark memory tuning starts by naming the region that ran out. Heap OOM means a task had too much data for its share of execution memory or user memory, so fix partitions and skew first. A container kill means native memory outside the heap, so raise overhead or budget Python. Size executors from the node down, compute per-task execution memory explicitly, and tune from measured maxima.