spark.sql.shuffle.partitions sets how many partitions Spark creates after a shuffle in DataFrame and SQL operations (joins, aggregations, window functions). The default is 200, whatever the data size. It does not apply to the starting partitions when reading files.
Why a fixed 200 is often wrong
- Small data: 200 MB shuffled into 200 partitions means 1 MB tasks. The scheduling overhead is bigger than the work, and you may write 200 tiny files.
- Large data: 2 TB shuffled into 200 partitions is 10 GB per partition. Tasks spill or run out of memory.
How to choose a value
A common rule is to aim for roughly 100 to 200 MB of shuffle data per partition. Estimate the shuffle size from the Spark UI (shuffle write of the previous stage), then divide.
shuffle size 400 GB / target 128 MB ~= 3,200 partitions
Also keep at least a few times the number of cores, so the cluster stays busy. With 400 cores, 200 partitions can never use more than half of them.
AQE makes this easier
With adaptive query execution on (it is on by default in Spark 3.2 and later), Spark looks at the real shuffle sizes at runtime and merges small adjacent partitions, aiming for an advisory size (spark.sql.adaptive.advisoryPartitionSizeInBytes, 64 MB by default). So you can set a deliberately high starting value, for example 2,000, and let Spark coalesce downwards for small jobs. AQE can only merge, never split, in this step. That is why starting high is safe and starting low is not.
Not the same as default parallelism
spark.default.parallelism only affects RDD operations such as reduceByKey or parallelize, and only when you do not give a partition count. DataFrame shuffles ignore it. Mixing them up is a common mistake in interviews.