Shuffle partitions
The default of 200 shuffle partitions, and how to size them for your data.
On this page
Show code in
Every code block on the page follows this.
You will learn
- What spark.sql.shuffle.partitions controls
- Why the default of 200 is wrong for most jobs
- How to size it from data volume
- How AQE changes the advice
Read first
- Partitions · 4 min read
- Narrow vs wide transformations · 3 min read
Comfortable with these? Read on.
spark.sql.shuffle.partitions (default 200) sets how many partitions a join or aggregation shuffles into. It is a fixed number for every shuffle in the session, so it is too many for small data and too few for large data, unless AQE adjusts it.What it controls
Every Exchange created by a DataFrame join, groupBy, distinct or window splits data into this many partitions, so the next stage runs this many tasks. It does not affect reading files (that is maxPartitionBytes) or explicit repartition(n) calls.
Why 200 is rarely right
| Shuffle size | With 200 partitions | Effect |
|---|---|---|
| 200 MB | 1 MB each | 200 tiny tasks, scheduling overhead dominates, and 200 small output files if written |
| 20 GB | 100 MB each | About right |
| 2 TB | 10 GB each | Huge tasks spill to disk, run for a long time, may fail with out-of-memory |
Sizing it
Divide the largest shuffle in the job by a target partition size. Read the shuffle size from the Spark UI (Stages tab, "Shuffle Write").
Rule of thumb
partitions ≈ largest shuffle size / 128-200 MB 2 TB / 200 MB ≈ 10,000 partitions then round to a multiple of the total cores, so the last wave of tasks is full
spark.conf.set("spark.sql.shuffle.partitions", "2000")
SET spark.sql.shuffle.partitions = 2000;
It is a session setting read when a query is planned, so you can change it between steps of a job when stages have very different sizes.
With AQE
Since Spark 3.2, AQE is on by default and coalesces small shuffle partitions after measuring them, toward spark.sql.adaptive.advisoryPartitionSizeInBytes (64 MB). That changes the advice:
- Set the initial number high enough for your largest shuffle. AQE only merges partitions; it does not split ordinary ones.
- AQE can choose the initial number itself from
spark.sql.adaptive.coalescePartitions.initialPartitionNumif set. - Small shuffles now come out right automatically, so the 200 MB example above becomes a handful of tasks.
Common mistakes
Leaving 200 for a multi-terabyte job
Setting it low and trusting AQE
Confusing it with spark.default.parallelism
Key takeaways
- spark.sql.shuffle.partitions (200) sets the partition count of every DataFrame shuffle.
- Size it from the largest shuffle: about 128 to 200 MB per partition.
- With AQE, start high and let it merge down.
- Streaming queries keep the number from their checkpoint.
Check yourself
3 questions1. What does spark.sql.shuffle.partitions affect?
Show the answer
The number of partitions after DataFrame shuffles. It sets the partition count for joins, aggregations and other exchanges.
2. A job shuffles 1 TB. Roughly how many shuffle partitions is reasonable?
Show the answer
5,000 to 8,000. 1 TB / 128-200 MB gives roughly 5,000 to 8,000.
3. With AQE on, should the initial number be low or high?
Show the answer
High enough for the largest shuffle; AQE merges small ones. Coalescing only merges, so start high.
Go deeper
Adaptive Query ExecutionSpark internals
PartitionsSpark internals
The small file problem
Primary sources: Performance tuning · Configuration