Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. What is spill?

PySpark · Memory, OOM & Tuning

What is spill?

Hardpyspark-56
spillshuffleexecution-memoryspark-ui

Question

What does "spill" mean in the Spark UI, and is it always bad?

Solution

Spill means Spark ran out of execution memory during a sort, aggregation or shuffle and had to write intermediate data to disk to continue. The job still finishes, it just does more disk I/O.

What you see in the UI

In the stage page, the summary shows two numbers:

  • Spill (memory): the size of the data in memory form that was spilled.
  • Spill (disk): the size it took on disk after serialization and compression.

The memory number is usually larger than the disk number, because the in-memory form is bigger. Do not be alarmed that the two differ.

Why it happens

Say a task has to sort or aggregate a 3 GB partition, but its share of execution memory is about 1 GB. Each executor's execution memory is divided among the tasks that run at the same time. Spark sorts what fits, writes it as a sorted run to disk, continues, and then merges the runs. Disk is far slower than memory, so tasks get slower.

Is it always bad

No. A small amount of spill on a few tasks costs little. Heavy spill across most tasks of a stage is a warning that the partitions are too large for the memory available. The stage will run, maybe 2 to 5 times slower than needed, and it increases the chance of failures.

Fixes, in order

  • More shuffle partitions, so each partition is smaller. For example moving from 200 to 800 on a 400 GB shuffle.
  • Let adaptive query execution coalesce and split partitions, and set sensible advisory partition sizes.
  • Reduce the data going in: filter and project columns before the shuffle.
  • Fix skew, if only a few tasks spill.
  • Give each task more memory, by lowering cores per executor or raising executor memory.

Spill is a symptom of "too much data per task". The fix is almost always to cut the data per task, and the memory change comes after that.

🎯 Put this concept into practice

Solidify this answer with real hands-on interview drills in the browser studio.

Open related drill →
PreviousNext