University · Basics · Lesson 3
How Spark memory works: executors, tasks and off-heap memory
ETL Compare staff · Published 18 September 2026 · 4 min read
In brief
Each Spark executor has a pool of execution memory that its running tasks share, so the memory one task gets depends on how many tasks run beside it. When a task needs more than its share, it spills to disk. Native accelerators add a second pool, off-heap memory, that you size yourself, which is why memory settings are part of every accelerator trial.
What happens when a task runs out of memory?
It spills: it writes part of its working data to local disk and reads it back later. The job still finishes, but the stage now pays for disk writes, disk reads and serialization. A task whose partition is much larger than the rest, because of data skew, is the usual victim, and it often becomes the straggler that holds up the whole stage.
Which Spark UI metrics show memory pressure?
- Peak execution memory: the most execution memory a task used
- Shuffle spill (memory): the size in memory of the data that was spilled
- Shuffle spill (disk): the size of that data once written to disk
- GC time: time the JVM spent in garbage collection, which often rises when memory is tight
All four appear per stage and per task in the Stages tab of the Spark web UI. Compare the median task with the maximum: a large gap on a slow stage points to skew rather than to a cluster that is too small.
What is off-heap memory and why do accelerators need it?
Off-heap memory is memory Spark manages outside the JVM heap. It is switched on with spark.memory.offHeap.enabled and sized with spark.memory.offHeap.size. Native engines keep their working data there, so the open-source accelerators ask you to size it: Gluten's Velox getting-started page uses 20g as its example, and Comet shares the off-heap pool with Spark.
Comet's tuning guide is direct about the limits. It notes that its memory accounting "isn't 100% accurate", and it caps native spill files per task with spark.comet.maxTempDirectorySize, which defaults to 100 GB. Choosing these values is part of the tuning and operating burden that our scores measure.
What do the usual fixes look like?
The general fixes are the same with or without an accelerator: give each task more memory (fewer cores per executor, or more memory per core), split the work into more and smaller partitions, fix skew so that no single task is oversized, and filter and project before wide operations. Adaptive query execution handles part of this automatically.
Accelerators approach the problem in different ways. The open-source native engines move working memory into the off-heap pool you size. DualBird states that its pipeline runs a handful of tasks at a time, each with a larger memory budget. Whatever the claim, the Spark UI metrics above are how you check it on your own job.
Related
Sources
- Web UI - Spark documentation · Fetched 27 Sep 2026
- Faster Spark Was the Goal. A Simpler Spark Came With It · Fetched 27 Sep 2026
- DualBird on skew and spill · Fetched 27 Sep 2026
- Apache Gluten: Velox backend getting started · Fetched 27 Sep 2026
- Apache DataFusion Comet tuning guide · Fetched 27 Sep 2026