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.