All posts

Spark memory management: every region, and the arithmetic that sizes it

An executor's memory is four regions inside the heap plus two outside it, and the sizes follow from three numbers. Here is the arithmetic, checked against a running executor, the elastic boundary read from the source, what the binary row format is measured to save, where PySpark's workers actually live, and how to tell which of the three layers killed your job.

27 min read Spark

TL;DR

  • An executor’s heap is split four ways: a fixed 300 MB reserved, a user region Spark does not manage, and one unified pool shared by execution and storage.
  • The sizes come from two fractions applied in order: (heap - 300MB) * spark.memory.fraction, then * spark.memory.storageFraction for storage’s protected share. On a 4 GB heap that is 2277.6 MiB unified, of which 1138.8 MiB is storage’s to keep.
  • Execution wins. The boundary between execution and storage is elastic and only storage gets evicted, so cached data is never guaranteed to still be there.
  • The container is killed for the total, not for the heap. Off-heap is additional rather than carved out of it, overhead is a third allocation on top, and Python workers are separate OS processes charged to that overhead: measured at roughly 55 MiB each before your code allocates anything.
  • Tungsten’s saving is not universal. Cached as a DataFrame rather than JVM objects, object-heavy rows measured 3.25x smaller, and a column of primitive ints measured no smaller at all.

Memory is where Spark tuning stops being folklore and becomes arithmetic. Given three configuration values you can compute every region size exactly, and given the region sizes most memory failures become predictable rather than mysterious.

The trouble is that the regions have confusing names, two of them overlap, one is invisible to Spark, and the number the cluster manager kills you for is none of the ones you configured. This post lays out each region, derives its size, checks the derivation against a running executor, and then goes through what happens when each one fills up.

Every number below was read from a live session or from the Spark source. Where I tried to demonstrate something and failed, the post says so.

Architecture: what the cluster manager allocates

Start outside the JVM, because this is the level at which your job gets killed. A container is three or four allocations, only one of which is the heap:

flowchart TB
  subgraph C["executor container, what YARN or Kubernetes reserves"]
    H["<b>JVM heap</b><br/>spark.executor.memory"]
    O["<b>off-heap</b><br/>spark.memory.offHeap.size<br/>(0 unless enabled)"]
    V["<b>memory overhead</b><br/>spark.executor.memoryOverhead<br/>JVM internals, native buffers"]
    P["<b>PySpark memory</b><br/>spark.executor.pyspark.memory<br/>(only for Python)"]
  end
  C --> K{"total > container limit?"}
  K -->|"yes"| X["container killed<br/>by the cluster manager"]
  K -->|"no"| R["runs"]

The point of that diagram is the arrow at the bottom. spark.executor.memory sizes the heap, and nothing else. The cluster manager enforces the sum, so a job that raises executor.memory to exactly the container limit is asking to be killed, because overhead still has to fit somewhere.

Config Default Covers
spark.executor.memory 1g The JVM heap, and only the heap
spark.executor.memoryOverhead unset JVM metaspace, thread stacks, native libraries, Python
spark.executor.memoryOverheadFactor 0.10, and 0.40 for non-JVM Kubernetes jobs The overhead as a fraction of heap, when overhead is not set directly
spark.executor.minMemoryOverhead 384m The floor under the factor calculation
spark.memory.offHeap.enabled false Whether Tungsten allocates outside the heap
spark.memory.offHeap.size 0 How much, when enabled

So overhead is max(executor.memory * 0.10, 384m) unless you set it directly, in which case the factor and the minimum are both ignored. That is worth knowing because the two derived knobs look like they still apply and they do not.

A naming change worth pinning down. spark.executor.minMemoryOverhead carries .version("4.0.0") in the config definitions, and it is absent from the 3.5.9 source, where the same 384 MiB floor was the hardcoded constant ResourceProfile.MEMORY_OVERHEAD_MIN_MIB = 384L. The number has not changed; it became configurable.

What are the four regions inside the heap?

Now inside the JVM. The heap is divided in a fixed order, and each division is a simple multiplication:

flowchart TB
  A["<b>JVM heap</b><br/>for example 4096 MiB"] --> B["<b>Reserved: 300 MB</b><br/>fixed, not configurable"]
  A --> C["<b>usable = heap - 300MB</b><br/>3796 MiB"]
  C --> D["<b>User memory</b><br/>usable x (1 - memory.fraction)<br/>1518 MiB, unmanaged"]
  C --> E["<b>Unified pool</b><br/>usable x memory.fraction<br/>2277.6 MiB"]
  E --> F["<b>Storage</b><br/>protected share:<br/>pool x storageFraction<br/>1138.8 MiB"]
  E --> G["<b>Execution</b><br/>the rest, and it can<br/>take storage's free space"]

Each region has a different job and a different failure mode:

Region Size Holds When it fills
Reserved 300 MB, fixed Nothing. It is a safety margin for Spark’s own internals Not applicable
User usable * (1 - memory.fraction) Your objects: UDF state, data structures in closures, anything Spark does not track OutOfMemoryError, with no spill and no warning
Storage pool * storageFraction protected Cached blocks, broadcast variables Blocks are evicted, then recomputed on next use
Execution The rest of the pool Shuffle buffers, sorts, hash tables for joins and aggregations Spills to disk

The user region is the one that surprises people. Spark does not track it, so it has no spill path and no eviction. If your UDF builds a large dictionary per task, it competes for the 40% of usable heap that Spark is deliberately not managing, and the failure is a plain OutOfMemoryError.

Does the arithmetic hold on a running executor?

The derivation above is easy to get wrong, so here it is checked against the MemoryManager in a live session with a 4 GB heap:

val mm = org.apache.spark.SparkEnv.get.memoryManager
def mb(b: Long) = f"${b / 1024.0 / 1024.0}%.1f MiB"

println("runtime maxMemory     = " + mb(Runtime.getRuntime.maxMemory))
println("maxOnHeapStorage      = " + mb(mm.maxOnHeapStorageMemory))
println("memory manager class  = " + mm.getClass.getSimpleName)

// the documented arithmetic, recomputed by hand
val heap     = Runtime.getRuntime.maxMemory
val usable   = heap - 300L * 1024 * 1024
val unified  = (usable * 0.6).toLong
println("usable  (heap - 300MB) = " + mb(usable))
println("unified (x 0.6)        = " + mb(unified))
println("storage region (x 0.5) = " + mb((unified * 0.5).toLong))
runtime maxMemory     = 4096.0 MiB
maxOnHeapStorage      = 2277.6 MiB
memory manager class  = UnifiedMemoryManager

usable  (heap - 300MB) = 3796.0 MiB
unified (x 0.6)        = 2277.6 MiB
storage region (x 0.5) = 1138.8 MiB

The hand calculation matches maxOnHeapStorageMemory exactly. And that raises the point most diagrams get wrong.

maxOnHeapStorageMemory is 2277.6 MiB, the whole unified pool, not the 1138.8 MiB storage region. Storage is not capped at its fraction. The fraction sets how much storage may keep under pressure; with execution idle, cached data can occupy the entire pool. The two numbers answer different questions:

  • spark.memory.storageFraction is the share of the pool that execution cannot take back.
  • maxOnHeapStorageMemory is the share that storage may grow into when nothing else wants it.

Read storageFraction as an eviction floor, not a size limit. That single reframing makes the next section obvious.

How does the boundary between execution and storage move?

Execution and storage share one pool, and the divide between them moves. The rules are asymmetric and they are short enough to read directly.

flowchart TB
  S1["storage using less<br/>than its share"] -->|"execution needs memory"| S2["execution takes<br/>the free space"]
  S3["storage borrowed past<br/>its protected share"] -->|"execution needs memory"| S4["blocks are <b>evicted</b><br/>until storage is back<br/>to its share"]
  E1["execution using more<br/>than half the pool"] -->|"storage needs memory"| E2["storage <b>waits</b>.<br/>Execution is never evicted,<br/>it only spills on its own terms"]

From UnifiedMemoryManager, this is what execution is allowed to reclaim:

        val memoryReclaimableFromStorage = math.max(
          storagePool.memoryFree,
          storagePool.poolSize - storageRegionSize)

The max of two things: storage’s free space, and however much storage has borrowed beyond its region. So execution can always take back what it lent, plus anything storage is not currently using.

And this is the cap on how large the execution pool may become:

    def computeMaxExecutionPoolSize(): Long = {
      maxMemory - math.min(storagePool.memoryUsed, storageRegionSize)
    }

min(storage used, storage region) is exactly the protected part. Execution may grow into everything except that.

There is no symmetric rule for storage. Nothing in the manager evicts execution memory, because a half-built hash table cannot be dropped and recomputed cheaply the way a cached block can. Execution releases memory by spilling to disk, and it decides when.

The practical consequence: cache() is a hint, not a guarantee. A cached DataFrame that fitted yesterday can be partly evicted today because a different stage in the same application wanted execution memory.

This is not how it always worked. Before 1.6 the split was static: a spark.shuffle.memoryFraction for execution and a separate spark.storage.memoryFraction for cache, each a fixed slice of the heap with no way to move memory between them. A job that cached nothing still paid for the storage region, and a job that cached heavily starved execution. SPARK-10000 replaced both with the single pool and the borrowing rule above, which is why those two settings no longer appear in the configuration reference.

What I could and could not demonstrate

Two claims from the section above are easy to state and harder to produce on demand, so here is what actually happened when I tried.

Storage borrowing past its region: confirmed. With a 1 GB driver heap, so a 434 MiB unified pool and a 217 MiB storage region, caching 7,000,000 rows put 306 MiB into storage across 8 blocks, comfortably past the region:

pool=434 MiB  region=217 MiB
after caching: storage=306 MiB blocks=8  beyond region=true

Eviction under execution pressure: not reproduced. I then ran sorts designed to demand execution memory, up to 30,000,000 rows over 8 concurrent tasks, and nothing was evicted:

after sort  : storage=306 MiB blocks=8  evicted=0
peakExecutionMemory (max task) = 0 MiB
memoryBytesSpilled  (total)    = 0 MiB
diskBytesSpilled    (total)    = 0 MiB

Zero spill is the explanation: there was never real pressure, so nothing needed reclaiming. The mechanism is in the source above and I am not going to claim I watched it fire when I did not. If you want to see eviction on your own cluster, the observable signals are the storage memory dropping in the executors tab and cached partitions falling below 100% in the storage tab.

A different way cache fails, and this one did reproduce. Caching more than the pool can hold with MEMORY_ONLY does not raise anything. It silently caches nothing:

# 20,000,000 rows into a 434 MiB pool
after caching: storage=0 MiB  blocks=0
cached count = 20000000     # still correct, recomputed from source

Zero blocks cached, and the count is still right because Spark recomputed the whole thing. A cache() call that achieves nothing while your job still produces correct answers is the worst kind of performance bug, because nothing fails. The storage tab showing a cached RDD at 0% is the tell, and MEMORY_AND_DISK is the usual fix.

When should you use off-heap memory?

Off-heap is Tungsten allocating with sun.misc.Unsafe outside the Java heap, so the data is not subject to garbage collection and is stored in Spark’s own binary format rather than as Java objects.

// started with --conf spark.memory.offHeap.enabled=true --conf spark.memory.offHeap.size=1g
println("maxOnHeapStorage   = " + mb(mm.maxOnHeapStorageMemory))
println("maxOffHeapStorage  = " + mb(mm.maxOffHeapStorageMemory))
println("tungstenMemoryMode = " + mm.tungstenMemoryMode)
maxOnHeapStorage   = 2277.6 MiB
maxOffHeapStorage  = 1024.0 MiB
tungstenMemoryMode = OFF_HEAP

Three things to take from that output:

  • Off-heap is additional. The on-heap pool is unchanged at 2277.6 MiB. Enabling 1 GB off-heap did not shrink the heap; it added a second pool, and the container now needs 1 GB more than before.
  • It has its own execution and storage split, governed by the same storageFraction, so the same elastic rules apply within it.
  • tungstenMemoryMode flips to OFF_HEAP, which is what actually changes where sorts and hash tables allocate.

The trade is GC pressure against manual accounting. Off-heap memory is invisible to the JVM, so a leak or an underestimate shows up as the cluster manager killing the container rather than as an OutOfMemoryError you can read a stack trace from.

What does the binary format actually buy?

Everything above treats the unified pool as a number of bytes. What those bytes hold is the other half of the story, because the same rows cost wildly different amounts depending on how they are represented.

Tungsten’s answer is to stop storing rows as JVM objects. UnsafeRow keeps a row in raw memory, and its source describes the layout exactly:

Each tuple has three parts: [null-tracking bit set] [values] [variable length portion]

One bit per field for nulls, then one 8-byte word per field. A fixed-width primitive such as int, long or double is stored directly in its word. A string or other variable-length value stores an offset and a length packed into that word, pointing into the third region. No object headers, no pointer chasing, and comparisons can run over raw bytes.

The folklore is that this makes cached data roughly four times smaller. It is worth measuring rather than repeating, so here is the same million rows cached six ways on Spark 4.1.3, local[2], four partitions, sizes read from sc.getRDDStorageInfo:

What is cached Form Bytes MiB
1M ints, sequential DataFrame 11,044 0.011
1M ints, random DataFrame 4,020,544 3.834
1M ints, random, compression off DataFrame 4,014,896 3.829
1M ints RDD[Int], MEMORY_ONLY 4,000,064 3.815
1M ints RDD[Int], MEMORY_ONLY_SER 5,000,000 4.768
1M Event(int, String, double) RDD[Event] 36,001,280 34.333
1M Event(int, String, double) DataFrame 11,065,712 10.553

Three things in that table are worth more than the headline.

For primitive ints, the RDD and the DataFrame cost the same. 3.815 MiB versus 3.834 MiB, with the DataFrame marginally larger. There is no fourfold saving because there is nothing to save: Scala specialises RDD[Int] into a primitive Array[Int], so there were never any object headers to remove. If you benchmark Tungsten with a column of integers you will conclude it does nothing.

The saving is real once rows hold objects. The same million rows as a case class with a String field cost 34.333 MiB as JVM objects and 10.553 MiB as a DataFrame, a 3.25x reduction. That is where the headers, the references and the per-object padding were, and that is what the binary format removes.

The largest number in the table has nothing to do with the row format. A sequential integer column caches in 11 KB, about 362 times smaller than the same column with random values, because the in-memory columnar cache compresses it. Turning spark.sql.inMemoryColumnarStorage.compressed off changed the random column by 0.1%, since random data does not compress. Sequential identifiers, low-cardinality categories and sorted columns can cache for almost nothing, and a benchmark built on spark.range will flatter every format it tests.

One more entry deserves a note: MEMORY_ONLY_SER made the integer RDD larger, 5,000,000 bytes against 4,000,064. Serialisation is a space win when it replaces object graphs, and a loss when it replaces a primitive array with a byte stream carrying per-element framing. It is not a free compaction.

The honest summary is that this is one shape of data on one version. The structural point holds regardless: the saving comes from removing per-object overhead, so it scales with how object-heavy your rows are, and it is invisible on primitives.

How do tasks share execution memory?

One more division, inside the execution pool. Several tasks run per executor and they compete, with a fairness rule in ExecutionMemoryPool:

      val maxPoolSize = computeMaxPoolSize()
      val maxMemoryPerTask = maxPoolSize / numActiveTasks
      val minMemoryPerTask = poolSize / (2 * numActiveTasks)
flowchart TB
  P["execution pool, N active tasks"] --> A["each task capped at <b>1/N</b><br/>of the pool"]
  P --> B["each task guaranteed at least <b>1/2N</b><br/>before it is made to wait"]
  B --> C["below 1/2N, the task blocks<br/>until another releases memory"]

So a task can never take more than 1/N of the pool, and it will not be blocked until it has at least 1/2N. N is the number of active tasks, not the core count, so the guarantee moves as tasks start and finish.

This is why spark.executor.cores is a memory setting as much as a parallelism setting. Doubling the cores per executor halves each task’s share of the same pool, which is the usual reason a job that worked with 4 cores per executor spills constantly with 8.

Why is there a 450 MiB floor?

Spark refuses to start if the heap cannot cover the reserved region with room to spare. From UnifiedMemoryManager:

  private val RESERVED_SYSTEM_MEMORY_BYTES = 300 * 1024 * 1024
  ...
    val minSystemMemory = (reservedMemory * 1.5).ceil.toLong

300 MB * 1.5 is 450 MiB, and there are two separate checks. A 400 MB driver:

[INVALID_DRIVER_MEMORY] System memory 419430400 must be at least 471859200.

and a 400 MB executor:

[INVALID_EXECUTOR_MEMORY] Executor memory 419430400 must be at least 471859200.

471859200 is exactly 450 MiB. This is a fail-fast check added deliberately, and it is a good thing: below that size the unified pool would be a few tens of megabytes and every job would spill pathologically instead of failing clearly.

What does PySpark add to the picture?

Everything so far is JVM memory. PySpark adds a second runtime that the JVM cannot see, and it is the single most common reason a Python job dies where the equivalent Scala job survives.

Python code does not run inside the executor JVM. Each executor starts separate Python worker processes and talks to them over sockets. You can make them identify themselves from inside a UDF:

import os
from pyspark.sql import SparkSession
from pyspark.sql.functions import udf
from pyspark.sql.types import StringType

spark = SparkSession.builder.appName("whoami").master("local[2]").getOrCreate()
print("driver python pid =", os.getpid(), " ppid =", os.getppid())

@udf(StringType())
def whoami(v):
    with open("/proc/self/status") as f:
        rss = [l.split()[1] for l in f if l.startswith("VmRSS")][0]
    return f"pid={os.getpid()} ppid={os.getppid()} rss={int(rss)//1024}MiB"

rows = (spark.range(0, 8, 1, 4)
        .selectExpr("CAST(id AS INT) AS v")
        .select(whoami("v").alias("worker"))
        .distinct().collect())
for r in sorted(set(x.worker for x in rows)):
    print(r)
spark.stop()

On 4.1.3 that prints a driver and two workers, one per core, each a distinct process with a different parent:

driver python pid = 1759  ppid = 1698
pid=1814 ppid=1810 rss=53MiB
pid=1814 ppid=1810 rss=56MiB
pid=1815 ppid=1810 rss=53MiB
pid=1815 ppid=1810 rss=56MiB

Two worker processes, 1814 and 1815, one per core. Each appears twice because its resident size grew between batches, which is the point: these are live processes with their own footprint, not a buffer inside the JVM. Their parent, 1810, is neither the driver nor its parent.

Roughly 55 MiB of resident memory per worker, before your code allocates anything, and none of it is in the heap you sized with spark.executor.memory. The plan names the boundary explicitly:

*(2) Project [pythonUDF0#11 AS p#10]
+- BatchEvalPython [plus(cast(id#0L as int))#9], [pythonUDF0#11]
   +- *(1) Range (0, 20000, step=1, splits=4)

BatchEvalPython is the operator that serialises rows out of the JVM, waits for Python, and reads them back. Note what surrounds it: the codegen markers *(1) and *(2) stop at its boundary, because a Python UDF is opaque to the whole-stage code generator. That is the CPU cost. The memory cost is the process it feeds.

Where does Python memory get charged?

To overhead, which is why the overhead default is different for Python. The factor is 0.10 normally, and Kubernetes uses NON_JVM_MEMORY_OVERHEAD_FACTOR, defined as 0.4d, for non-JVM applications. The description on spark.kubernetes.memoryOverheadFactor states it plainly: the default is “0.10 and 0.40 for non-JVM jobs”.

That 0.40 is worth reading as a warning rather than a gift. It is the project’s estimate that a Python executor needs four times the non-heap headroom of a JVM one, and it applies on Kubernetes; on YARN you may be running a Python job with a JVM-sized overhead unless you set it yourself.

Config Default on 4.1.3 What it does
spark.executor.pyspark.memory unset Reserves a dedicated pool for Python workers, accounted separately from overhead
spark.executor.memoryOverheadFactor 0.10 Overhead as a fraction of heap, when overhead is not set directly
spark.sql.execution.arrow.maxRecordsPerBatch 10000 Rows per Arrow batch moved between JVM and Python

spark.executor.pyspark.memory is unset by default, which means Python competes with JVM internals, native libraries and Arrow buffers for one undifferentiated overhead budget. Setting it does not create memory; it partitions the overhead so that a Python worker leak fails as a Python worker limit rather than as an unexplained container kill.

spark.sql.execution.arrow.maxRecordsPerBatch is the knob people find last and should find first. Arrow batches are allocated outside the heap, so the batch size multiplies directly into the overhead region. Too small wastes the vectorised path; too large turns a working job into a container kill with no JVM error. It is a memory setting wearing a throughput setting’s name.

Three failures that all look like “PySpark ran out of memory”

Symptom Where it actually failed What to change
java.lang.OutOfMemoryError in the executor log JVM heap, execution or storage The heap, cores per executor, or the query shape
Python traceback ending in MemoryError, or the worker dying mid-stage The Python worker process spark.executor.pyspark.memory, or stop materialising whole partitions in pandas
Exit 137 with no JVM stack trace at all The container total, usually overhead spark.executor.memoryOverhead, or a smaller Arrow batch

The middle row is the one that has no Spark-side signal. A Python UDF that calls toPandas() on a partition, or builds a large dictionary per task, allocates in a process Spark does not measure and no memory metric will show it.

What about the driver?

The driver gets far less attention than the executor and fails in its own ways. It runs no tasks, holds no shuffle data, and caches nothing, so it has no unified pool worth reasoning about. Its pressure comes from three places: results pulled back by an action, broadcast variables on their way out, and the scheduler’s own bookkeeping for a large plan.

The defaults, read from the running build:

Config Default on 4.1.3 Covers
spark.driver.memory 1g The driver JVM heap
spark.driver.memoryOverheadFactor 0.1 Overhead as a fraction of driver heap
spark.driver.memoryOverhead unset, derived from the factor JVM internals, and in client mode the Python or R process
spark.driver.maxResultSize 1g Total serialised results one action may return

spark.driver.maxResultSize is the one to understand, because it is a guard rather than a limit on memory. When an action would return more than a gigabyte of serialised results, Spark aborts the job with an error naming the setting, instead of letting the driver collect its way into an OutOfMemoryError. The error is the system working. Raising it is reasonable when you know the result is bounded and merely larger than a gigabyte. Raising it because a collect() on an unbounded DataFrame failed is removing the airbag.

Three driver failures cover most of what happens in practice:

  • collect() on something unbounded. A groupBy with no aggregation, or a filter that turns out not to be selective. Write to storage and read it back instead.
  • A broadcast that is not small. The driver assembles the broadcast relation before sending it. A “small” dimension table that grew past the broadcast threshold fails on the driver, not the executor.
  • local[*], where the driver is also the executor. Every executor-side concern in this post applies to the driver process, and spark.executor.memory is ignored. Size spark.driver.memory.

How do you tell which layer failed?

Memory failures in Spark are not one failure. They come from three different layers, with three different observers, and the first job in any incident is to work out which one you are looking at.

Layer Who reports it Signature
Spark’s own accounting Spark Spill in the logs, blocks evicted, a stage that slows without failing
The JVM process The JVM java.lang.OutOfMemoryError with a stack trace
The operating system, via the cluster manager YARN or Kubernetes Exit 137, OOMKilled, “X GB of Y GB physical memory used”, no stack trace

Only the first is visible in Spark’s memory metrics. The third is invisible to Spark entirely, because the process is killed from outside.

flowchart TB
  S["memory incident"] --> Q1{"is there a JVM<br/>stack trace?"}
  Q1 -->|"no, exit 137<br/>or OOMKilled"| C1["the container total was exceeded"]
  C1 --> Q2{"PySpark, Arrow,<br/>or native libraries?"}
  Q2 -->|"yes"| F1["overhead, or<br/>spark.executor.pyspark.memory"]
  Q2 -->|"no"| F2["off-heap size,<br/>or shrink the heap"]
  Q1 -->|"yes, OutOfMemoryError"| Q3{"is there spill<br/>in the logs?"}
  Q3 -->|"yes"| F3["execution pressure:<br/>fewer cores, or raise memory.fraction"]
  Q3 -->|"no"| F4["user memory:<br/>your own objects, no spill path"]
  Q1 -->|"no error,<br/>just slow"| Q4{"cached fraction<br/>below 100%?"}
  Q4 -->|"yes"| F5["storage evicted:<br/>raise storageFraction or use MEMORY_AND_DISK"]
  Q4 -->|"no"| F6["not a memory problem"]

Two practical notes on reading the evidence.

The Spark UI does not show the process total. The storage and executor tabs report what the UnifiedMemoryManager tracks, which excludes the user region, native allocations, Arrow buffers and every Python worker. A container can be at its limit while the UI shows a comfortable heap. The gap between “Spark-tracked memory” and “what the kernel counts” is exactly where container kills live.

Process-level metrics are available and off by default. spark.executor.processTreeMetrics.enabled defaults to false on 4.1.3. Turning it on adds process-tree memory to the executor metrics, which is the closest Spark gets to reporting the number the cluster manager enforces. It costs some sampling overhead, and it is worth it on any Python workload you are trying to size.

There is also a limit to what tuning can fix, and it is worth saying plainly. Every knob in this post redistributes a fixed container budget. If a single task genuinely needs more memory than one executor can be given, because a join explodes, a partition is badly skewed, or the table layout forces enormous state into memory at once, no fraction will save it. At that point the fix is in the data: partitioning, file layout, join strategy, or the query itself. Recognising that boundary early is worth more than another round of spark.memory.fraction.

What should you tune, and in what order?

Symptom First move Why
Container killed, exit 137, “memory overhead exceeded” Raise spark.executor.memoryOverhead The heap was never the problem; the sum exceeded the container
Heavy spill during joins and aggregations Fewer cores per executor, or raise spark.memory.fraction Each task gets 1/N of the pool, so N is the lever
Cached data keeps being recomputed Raise spark.memory.storageFraction, or use MEMORY_AND_DISK The fraction is the eviction floor
OutOfMemoryError with no spill in the logs Look at user memory, not the pool Spark does not manage or spill your own objects
GC pauses dominating Enable off-heap, or shrink the heap and add executors Off-heap data is not scanned by the collector
PySpark workers killed spark.executor.pyspark.memory, and raise overhead Python lives outside the JVM entirely

The ordering matters more than any individual row. Overhead problems masquerade as heap problems, and the first question for any kill is always whether the JVM died or the container was terminated from outside.

Common misconceptions

“spark.memory.storageFraction caps how much I can cache.” It does not. It sets the share of the unified pool that execution is not allowed to evict. maxOnHeapStorageMemory on the 4 GB heap above reports 2277.6 MiB, the whole pool, not the 1138.8 MiB storage region, because storage can grow into the entire pool when execution is idle. Read the fraction as an eviction floor.

“Enabling off-heap memory takes it out of the heap.” It is additional. Turning on 1 GB of off-heap left the on-heap pool unchanged at 2277.6 MiB and added a second pool of 1024.0 MiB. The container now needs a gigabyte more than before, which is how enabling off-heap gets a job killed by the cluster manager without a single change to spark.executor.memory.

“The container was killed, so the heap was too small.” The cluster manager enforces heap + off-heap + overhead + PySpark, and only the first of those is spark.executor.memory. A heap sitting at 60% with Python and native buffers pushing the container over its limit is killed with no JVM OutOfMemoryError at all, usually as exit code 137.

“An OutOfMemoryError means I should raise spark.memory.fraction.” Not if the memory went to your own objects. The user region, the part Spark does not manage, has no spill path and no eviction, so a UDF building a large structure per task fails outright. Raising the managed fraction makes that region smaller.

Production notes

  • Size the container, not the heap. heap + off-heap + overhead + PySpark is the number enforced, so leave headroom rather than setting executor.memory to the limit.
  • Setting spark.executor.memoryOverhead directly disables both derived knobs. The factor and the minimum are ignored once it is explicit.
  • Treat storageFraction as an eviction floor. It does not cap how much cache can grow, only how much survives pressure.
  • Never assume a cached DataFrame is resident. Check the storage tab for the cached fraction before attributing a slow stage to something else.
  • Use MEMORY_AND_DISK unless you have measured that the data fits. MEMORY_ONLY silently caches nothing when it does not, and the job stays correct while getting slower.
  • Treat spark.executor.cores as a memory knob. It divides the execution pool 1/N, so raising it can cause spill with no other change.
  • Account for off-heap twice. Once in offHeap.size and once in the container total. It is additional to the heap, not carved from it.
  • Do not tune spark.memory.fraction first. It trades user memory against managed memory, and the common problems are overhead and core count.

Frequently asked questions

Why is my container killed when the heap looks fine? Because the cluster manager enforces the total, and overhead is not part of spark.executor.memory. A heap sitting at 60% with off-heap and Python pushing the container over its limit gets killed with no JVM OutOfMemoryError at all, usually as exit code 137.

What is actually in the 300 MB reserved region? Nothing of yours. It is a margin so that Spark’s own internal structures cannot be squeezed out by execution and storage. It is not configurable outside tests, and the 450 MiB startup floor exists to keep it meaningful.

Does spark.memory.storageFraction limit how much I can cache? No, and this is the most common misreading. It sets the portion of the unified pool that execution cannot evict. With execution idle, cache can occupy the whole pool, which is why maxOnHeapStorageMemory reports the full pool size.

Can storage evict execution? No. The asymmetry is deliberate: a cached block can be dropped and recomputed, while a partially built hash table cannot. Execution gives memory back by spilling, on its own schedule.

Is off-heap memory faster? It is not automatically faster. It removes garbage-collection pressure and stores data compactly, which helps on large heaps, and it costs you the JVM’s memory accounting. A leak becomes a killed container instead of a stack trace.

Where does PySpark memory come from? Outside the JVM, from the overhead allocation, unless you set spark.executor.pyspark.memory to reserve it explicitly. A Python UDF holding a large object is not visible in any Spark memory metric.

Why did my job start spilling after I gave executors more cores? Because each task is capped at 1/N of the execution pool where N is the active task count. More cores means more concurrent tasks and a smaller share each, from the same pool.

Conclusion

The useful way to hold all of this is as one sum and one asymmetry.

The sum is the container: heap, off-heap, overhead and Python, of which spark.executor.memory is only the first. Most memory incidents that look mysterious are the cluster manager enforcing that sum while you were watching the heap.

The asymmetry is inside the unified pool: execution can take memory from storage, and storage can never take it from execution. Every consequence follows from that. Cached data is evictable, so caching is a hint. storageFraction is a floor rather than a ceiling, so it does not limit cache growth. Execution handles its own shortfall by spilling, which is why spill is a tuning signal rather than an error. Once those two ideas are in place the configuration stops being a list of fractions to try and becomes arithmetic you can do before changing anything: compute the pool from the heap, divide by the concurrent task count, and compare it with what one task actually needs. That prediction beats trial and error.

References

Trademarks

Apache Spark, Apache Hadoop and Apache are either registered trademarks or trademarks of The Apache Software Foundation in the United States and other countries.

Found this useful?

These posts and tools are free. If one saved you an afternoon, you can buy me a coffee.

Buy me a coffee