Skip to content
Good engineers know 3 min read · spark.sql.shuffle.partitions

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

Comfortable with these? Read on.

TL;DR 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 sizeWith 200 partitionsEffect
200 MB1 MB each200 tiny tasks, scheduling overhead dominates, and 200 small output files if written
20 GB100 MB eachAbout right
2 TB10 GB eachHuge 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
PySparkSpark SQL · Setting it
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.initialPartitionNum if set.
  • Small shuffles now come out right automatically, so the 200 MB example above becomes a handful of tasks.
Note: Structured Streaming ignores later changes: the number of shuffle partitions is stored in the checkpoint when a streaming query first starts, and cannot be changed without a new checkpoint.

Common mistakes

Leaving 200 for a multi-terabyte job

Gigantic tasks, spills and OOMs.

Setting it low and trusting AQE

AQE merges; it does not split normal partitions.

Confusing it with spark.default.parallelism

That one applies to RDD operations, not DataFrame shuffles.

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 questions

1. 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

Primary sources: Performance tuning · Configuration