ETL CompareETL Compare

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.

Where does a task's memory come from?

An executor is a JVM process that runs several tasks at once, one per core it is given. The tasks on one executor draw on the same execution memory, which Spark uses for sorts, joins, aggregations and shuffle buffers. Spark does not hand each task a fixed slice up front: the share moves as tasks start and finish.

DualBird's engineering article 'A Simpler Spark' sets out the arithmetic: with N active tasks, each task is guaranteed roughly 1/(2N) of the pool before it is forced to spill and is capped near 1/N. Its example of 4 concurrent tasks sharing 16 GiB leaves each task about 1 to 2 GiB. The point holds whichever product you run: adding cores to an executor without adding memory shrinks what each task can hold.

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