What does setting spark.memory.offHeap.enabled to true change about a Spark executor's memory?
answer
- a second pool, not a bigger heap
- the collector never looks at it
- two settings, useless without both
- Tungsten's binary pages move house
basics
~20 sIt gives Spark a second memory pool of spark.memory.offHeap.size bytes, allocated outside the JVM heap and used by Tungsten for execution and by the OFF_HEAP storage level. That memory is not garbage collected, so large working sets stop driving GC pauses.
solid answer
~50 sWith `spark.memory.offHeap.enabled=true` and a positive `spark.memory.offHeap.size`, Spark allocates a native memory region and adds it to the memory manager alongside the on-heap unified pool. Tungsten's operators — sort, hash aggregate, shuffle — can then hold their binary pages there instead of on the heap, and `StorageLevel.OFF_HEAP` caches blocks there serialized. The payoff is garbage collection: the JVM never scans these bytes, so a job with multi-gigabyte working sets stops paying long GC pauses. The costs are real too. The size is fixed at startup rather than elastic, it must be included in the container request on YARN and Kubernetes (Spark accounts for it, but your total container grows), and it is harder to observe — an overrun surfaces as a container kill rather than a familiar `OutOfMemoryError` with a heap dump. Setting the size without enabling the flag does nothing at all.
code
properties · 6 lines# both are required; either alone is useless
spark.memory.offHeap.enabled true
spark.memory.offHeap.size 4g
# container = executor.memory + memoryOverhead + 4g
spark.executor.memory 8g
spark.executor.memoryOverhead 2ggo deeper
Know only that Spark can hold some working data outside the Java heap and that this is an advanced tuning option, not something you set on an ordinary job.
Name both spark.memory.offHeap.enabled and spark.memory.offHeap.size, note that they must be set together, and explain that Tungsten's binary row format is what makes off-heap storage possible.
Justify the change with evidence from GC time in the task metrics, and be honest about the costs: a fixed non-elastic size, a larger container request, and worse diagnosability when it runs out.
Decide whether off-heap belongs in a platform's default executor shape at all, weighing GC stability on large jobs against container budget, queue density, and the support cost of failures that arrive without a heap dump.
## What "off-heap" means here Spark's Tungsten engine already stores rows in a compact binary format rather than as JVM objects. Because those rows are just bytes with a known layout, they do not need to live on the Java heap at all. Turning on off-heap memory tells Spark to allocate native memory directly and place those binary pages there. The JVM garbage collector never walks them. This is distinct from — and additional to — the on-heap unified pool. Enabling off-heap does not shrink or replace `spark.memory.fraction`; you now have two pools, and the memory manager tracks both. ## The two settings, and the trap - `spark.memory.offHeap.enabled` — defaults to `false`. - `spark.memory.offHeap.size` — defaults to `0`; must be positive when the feature is enabled. The trap is that they must be set together. Setting the size while leaving the flag off silently does nothing, which produces the classic "I configured off-heap memory and saw no change" report. Enabling the flag with a zero size fails at startup. ## Who uses the off-heap pool **Execution**: Tungsten's sort, hash-aggregate and shuffle operators take their pages from whichever pool is configured. Their spill behaviour is unchanged — running out of off-heap execution memory triggers the same flush to local disk it would on-heap. **Storage**: only if you explicitly persist with `StorageLevel.OFF_HEAP`, which stores blocks serialized in the off-heap region. Ordinary `MEMORY_AND_DISK` caching still uses the heap. ## Why anyone bothers: garbage collection A Spark executor with a large heap and a large live working set is a hard case for the JVM collector. Every full collection must trace the live set, and the live set here is exactly the multi-gigabyte hash table or sort buffer the job depends on. High "GC Time" in the task metrics — a substantial share of task duration — is the symptom. Moving that data off-heap removes it from the collector's reach entirely: the pauses shrink, and tail latency on long stages stabilises. A reasonable pattern is a smaller heap with a meaningful off-heap pool, rather than one very large heap. The container total stays comparable while the collector's job gets much easier. ## The costs - **It is fixed, not elastic.** The on-heap unified pool lets execution and storage borrow from each other dynamically. The off-heap size is a number you commit to at startup; if it is too small the operators spill more, and if it is too large you have wasted container budget the heap could have used. - **It must be in the container budget.** On YARN and Kubernetes, Spark includes `spark.memory.offHeap.size` in the resource request alongside `spark.executor.memory` and `spark.executor.memoryOverhead`, so enabling it makes containers larger — and a queue with fixed capacity will schedule fewer of them. - **Observability gets worse.** Heap problems produce an `OutOfMemoryError`, a stack trace, and a heap dump you can analyse. Native allocation problems produce a container kill or a process death with far less to go on. - **It does not fix leaks in user code.** Objects your UDFs allocate stay on the heap in user memory no matter what this setting says. ## When it is actually worth turning on Enable it when you have *measured* GC as a material share of stage time on large-working-set jobs — big sort-merge joins, wide aggregations, heavy shuffles — and when your executors are already large enough that a bigger heap is not the answer. Leave it off for ordinary jobs. It is a tuning step you take with evidence, not a default to sprinkle on a submit script. It is also worth being explicit about what it is *not*. Off-heap memory is not durable and not shared: it dies with the executor, exactly like the heap. It is not a cache tier that other applications can read. And it is not the same thing as the memory overhead allowance — the overhead covers uncontrolled native consumers such as network buffers and Python workers, whereas the off-heap pool is memory Spark itself manages and accounts for precisely.
- Does enabling off-heap memory remove the need for spark.executor.memoryOverhead?No. They cover different things. The off-heap pool is memory Spark itself allocates and accounts for precisely, and it is added to the container request as its own term. The overhead allowance covers consumers Spark does not manage — Netty shuffle buffers, metaspace and thread stacks, native codecs, Python workers. Enabling off-heap execution memory does not reduce any of those.
- How would you know off-heap memory is worth enabling for a particular job?Measure GC first. In the Spark UI's task metrics, compare GC Time against task Duration on the stages that dominate runtime. If garbage collection is consuming a material share of that time on stages with large sort or aggregation working sets, moving Tungsten's pages off-heap is likely to help. If GC time is small, the problem is elsewhere and this setting will change nothing.
- What happens if a Spark job needs more off-heap execution memory than spark.memory.offHeap.size allows?The same thing that happens on-heap: the spillable operators flush their pages to the executor's local disk and continue. The pool size is not a correctness boundary. What you lose is speed, and because the off-heap size is fixed at startup rather than borrowed dynamically from a storage region, an undersized value produces steadier spilling than the on-heap arrangement would.
saying these in an interview costs you the question
- Sets offHeap.size without enabling the flag and expects an effect
- Thinks off-heap memory survives an executor restart
- Confuses the off-heap pool with spark.executor.memoryOverhead
- Believes it eliminates spilling rather than relocating the pages
- Enables it by default without measuring GC time first