Skip to content
Reliable Data Engineering
Lesson
Open in the interactive app

Spark Memory Architecture and OOM Debugging

“The job died with an OutOfMemoryError. What do you do?” is one of the most common senior Spark questions, and one of the most commonly fumbled. “Add more memory” is the junior answer. The senior answer identifies which process ran out (driver or executor), which memory region (heap, off-heap, overhead, Python), why (skew, too-large partitions, collect, broadcast, caching, UDFs), and fixes the cause. This module gives you the mental model and the debugging path.


1. The memory layout of an executor

flowchart TB
    subgraph CONT["Container / pod limit = heap + overhead (+ off-heap + PySpark memory)"]
        subgraph HEAP["JVM heap: spark.executor.memory"]
            RES["Reserved: 300 MB<br/>Spark internals"]
            subgraph UNI["Unified memory = (heap − 300 MB) × spark.memory.fraction (0.6)"]
                EXE["Execution<br/>shuffle buffers, sort, hash aggregation,<br/>join hash tables"]
                STO["Storage<br/>cached blocks, broadcast variables,<br/>unrolled blocks"]
            end
            USR["User memory = (heap − 300 MB) × 0.4<br/>your objects, UDF data structures,<br/>Spark internal metadata"]
        end
        OFF["Off-heap (spark.memory.offHeap.size)<br/>optional Tungsten off-heap pool"]
        OVH["memoryOverhead: max(384 MB, 10%)<br/>JVM native memory, thread stacks, NIO/Netty buffers,<br/>metaspace, Python workers*, Arrow buffers*"]
        PYM["spark.executor.pyspark.memory (optional)<br/>caps Python worker memory"]
    end

* If spark.executor.pyspark.memory isn’t set, Python workers draw from the overhead allowance.

Worked example: spark.executor.memory=16g, 4 cores.

The borrowing rules (the asymmetry interviewers ask about)

The driver

The driver has the same JVM structure (spark.driver.memory, spark.driver.memoryOverhead), but it holds very different things:

spark.driver.maxResultSize (default 1 GB) caps the total serialized size of results collected from tasks, so a failing job is better than a dead driver.


2. Every kind of out-of-memory failure, decoded

Error / symptomWhich memoryTypical causesFirst fixes
java.lang.OutOfMemoryError: Java heap space (executor)Executor heap (execution/user)Huge partition (skew or too few partitions), giant rows, collect_list on big groups, wide explode, big UDF data structuresFind the task (skew?), more shuffle partitions, fix skew, reduce row width, avoid unbounded aggregates
OutOfMemoryError: Java heap space (driver)Driver heapcollect()/toPandas() of big data, broadcasting a large table, too many tasks/partitions (task metadata), huge plans from loops of withColumnDon’t collect, limit or aggregate first, lower broadcast thresholds, coalesce tiny partitions, build plans with select instead of loops
GC overhead limit exceededHeap, thrashingHeap nearly full, GC reclaims almost nothing; too many objects (RDD/Python-object-heavy code, deserialized caches)More memory per task, serialized caching, DataFrames instead of RDDs, G1GC tuning, fewer cores per executor
Container killed by YARN for exceeding memory limits / Kubernetes OOMKilled (exit 137)Total container: usually overhead or PythonPySpark workers, pandas UDFs and Arrow batches, off-heap/native libraries, Netty buffers, too many cores sharing overheadRaise memoryOverhead (or set pyspark.memory), smaller Arrow batches, fewer cores per executor, native functions instead of UDFs
OutOfMemoryError: Direct buffer memoryOff-heap NIO buffersLarge shuffle fetches / network buffers, some connectorsRaise overhead / -XX:MaxDirectMemorySize, reduce spark.reducer.maxSizeInFlight, more partitions
OutOfMemoryError: MetaspaceMetaspace (class metadata, in overhead)Many dynamically generated classes (codegen on huge plans), class loader leaks in long-running sessionsSimplify plans, restart long-lived sessions/notebooks, raise metaspace
SparkOutOfMemoryError: Unable to acquire N bytes of memoryExecution pool could not grant memory for an operatorOne task’s operator needs more than it can get (skewed join/aggregation, huge sort buffer)Same as heap OOM: skew, partition size, cores per executor
FetchFailedException / ExecutorLostFailure loopsIndirect: an executor died (often OOM or preemption)The executor holding shuffle output was killedFind why it died (usually one of the above), external shuffle service/decommissioning
spark.driver.maxResultSize exceededDriver protectionCollecting too muchWrite to storage instead of collecting

3. The debugging path

flowchart TD
    A[Job failed with OOM / container killed] --> B{Driver or executor?}
    B -->|"driver log, 'Driver stacktrace', app dies at once"| D1[Driver path]
    B -->|"task failures, ExecutorLostFailure, exit 137 on executors"| E1[Executor path]
    D1 --> D2{"collect / toPandas / take?<br/>broadcast of a big table?<br/>millions of tiny partitions?<br/>huge plan?"}
    D2 --> D3[Remove the collect, limit/aggregate,<br/>fix broadcast stats/hints, coalesce, simplify plan,<br/>then raise driver memory if still needed]
    E1 --> E2{Which stage and task?<br/>Stages tab: max vs median}
    E2 -->|"one or few tasks huge"| SK[Skew: see the skew guide]
    E2 -->|"all tasks large"| E3{Spill and GC metrics}
    E3 -->|"heavy spill everywhere"| P[Partitions too big:<br/>more shuffle partitions / AQE advisory size]
    E3 -->|"high GC time"| G[Too many objects per task:<br/>fewer cores per executor, serialized cache,<br/>DataFrames over RDDs, G1GC]
    E3 -->|"container killed, heap looks fine"| O[Overhead / Python memory:<br/>raise overhead or pyspark.memory,<br/>smaller Arrow batches, avoid UDFs]

Step 1: Driver or executor?

Step 2: Which stage and task?

In the Stages tab, open the failing stage and look at the Summary Metrics for completed tasks: duration, shuffle read size, records, spill, GC time (min / 25th / median / 75th / max).

Step 3: Spill

Spill (Memory) vs Spill (Disk): spill means operators already ran out of execution memory and wrote to disk. A ratio of memory to disk of roughly 3-10× is normal (deserialized vs compressed). Heavy spill on every task is the precursor to OOM: increase partitions or memory per task.

Step 4: GC time

GC time above ~10% of task time signals memory pressure; above 20-30% is a problem. Patterns:

Step 5: GC logs (when it isn’t obvious)

Add -Xlog:gc*:file=… (JDK 11+) or -verbose:gc via spark.executor.extraJavaOptions. Look for back-to-back full GCs that reclaim little (heap genuinely too small for the working set), long pauses, or humongous allocations with G1 (very large arrays such as big join hash tables or large records).

Step 6: Fix the cause, then size

Only after removing causes (skew, collect, oversized broadcast, UDF memory) should you resize: more memory per task (fewer cores or more memory), more overhead for PySpark, more partitions.


4. GC and serialization tuning


5. Observability: where to look

PlaceWhat it tells you
Executors tabStorage memory used per executor, task time vs GC time, failed tasks, shuffle read/write, peak JVM/off-heap/Python memory metrics (Spark 3.x executor metrics)
Stages tab → Summary MetricsDistribution of duration, shuffle read, spill, GC: skew vs uniform pressure
Storage tabWhat’s cached, fraction cached, memory vs disk size: over-caching
SQL tabPer-operator metrics (peak memory, spill size, rows) in the query DAG
Environment tabThe memory settings actually in effect (catches config that didn’t apply)
REST API / metrics sinks/api/v1/applications/<id>/executors for automation; the Prometheus servlet and executor metrics for Grafana dashboards and alerts (e.g. alert on GC ratio > 20%)
Cluster logsContainer exit codes: 137 = killed (OOM killer / limit), 143 = terminated

6. Memory anti-patterns

  1. collect() / toPandas() on unbounded data. Aggregate or write out instead.
  2. Broadcasting a table that’s large after decompression (or is large because statistics are stale).
  3. collect_list/collect_set over unbounded groups (one key → millions of elements in one row).
  4. Python UDFs on large columns without budgeting memory overhead.
  5. Building plans in a loop of withColumn (hundreds of nested projections → driver memory and planning time); use one select with all expressions.
  6. Fat executors (all cores, huge heap) → long GC pauses and shared-overhead OOMs.
  7. Caching everything “just in case”: evicts useful data and adds GC pressure.
  8. Too few shuffle partitions for the data (each task gets GBs).
  9. Ignoring skew and “fixing” it with bigger executors, which only moves the cliff.

7. Configuration reference

SettingDefaultWhat it controls
spark.executor.memory / spark.driver.memory1gJVM heap
spark.executor.memoryOverheadmax(384 MB, 10% of heap)Non-heap container allowance
spark.executor.pyspark.memoryunsetSeparate cap for Python workers
spark.memory.fraction0.6Unified (execution + storage) share of usable heap
spark.memory.storageFraction0.5Storage share of unified memory protected from eviction
spark.memory.offHeap.enabled / .sizefalse / 0Tungsten off-heap pool (also counts toward container size)
spark.driver.maxResultSize1gCap on collected results
spark.sql.shuffle.partitions200Partitions after shuffles (size per task)
spark.sql.adaptive.advisoryPartitionSizeInBytes64 MBAQE target partition size
spark.sql.autoBroadcastJoinThreshold10 MBMax estimated size to broadcast
spark.sql.execution.arrow.maxRecordsPerBatch10000Arrow batch size for pandas UDFs
spark.reducer.maxSizeInFlight48 MBShuffle fetch buffer per reduce task

Interview questions

Explain Spark's unified memory model and why execution can evict storage but not the reverse.

After 300 MB reserved, spark.memory.fraction (0.6) of the heap is a unified pool shared by execution (shuffles, sorts, aggregations, join hash tables) and storage (cache, broadcast). Either side can borrow the other’s free memory. Execution can evict cached blocks down to the protected storageFraction, but storage can’t evict execution memory, because tasks need execution memory to make progress while cached data can be recomputed or re-read. The remaining 0.4 is user memory for user objects and internal metadata.

A PySpark job fails with "Container killed by YARN for exceeding memory limits", but heap usage looks fine. Why?

YARN enforces the total container size, which includes memory outside the JVM heap: Python worker processes, Arrow buffers for pandas UDFs, native/NIO buffers and metaspace. These come out of memoryOverhead (10% by default), which is too small for Python-heavy workloads, especially with many cores (one Python worker per concurrent task). Fix: raise memoryOverhead or set spark.executor.pyspark.memory, reduce arrow.maxRecordsPerBatch, use fewer cores per executor, and replace Python UDFs with native functions.

How do you tell a driver OOM from an executor OOM, and what are the usual causes of each?

A driver OOM kills the whole application, appears in the driver log and usually follows an action that brings data to the driver (collect, toPandas), a broadcast (relations are collected on the driver first), or planning of a huge plan or millions of partitions. An executor OOM shows as failed and retried tasks with ExecutorLostFailure or exit 137 in executor logs, usually caused by skewed or oversized partitions, unbounded per-key aggregates, or UDF memory.

What does high GC time tell you and how do you fix it?

The heap is under pressure with many live objects: the JVM spends its time collecting instead of computing. If it’s high on all tasks: too many concurrent tasks per heap (reduce cores per executor or raise memory), object-heavy code (move from RDDs/Python objects to DataFrames), or deserialized caches (use serialized/columnar caching, unpersist). If only a few tasks: skew. Then tune G1GC (earlier concurrent marking, region size) if needed.

Why can increasing executor memory fail to fix an OOM?

If the cause is skew, one key’s data still lands in one task and keeps growing with the data; if it’s the driver (collect/broadcast), executor memory is irrelevant; if it’s overhead/Python memory, a larger heap doesn’t help (and can shrink the overhead headroom within a fixed container). Identify the process, the memory region and the root cause first.