How do you choose a value for spark.sql.shuffle.partitions on a large job?
answer
- the default was picked without seeing your data
- start from bytes, finish with cores
- a few hundred megabytes per task is comfortable
- too low kills the job, too high just wastes
basics
~20 sSize it from the shuffled data volume, targeting post-shuffle partitions of a few hundred megabytes each, and keep the count a comfortable multiple of total executor cores. The default of 200 is a fixed number that ignores your data size entirely.
solid answer
~40 s`spark.sql.shuffle.partitions` sets how many partitions every SQL/DataFrame exchange produces, and its default is **200** — a constant, chosen with no knowledge of your data. Size it from the shuffled bytes: divide the stage's shuffle-write volume by a target partition size in the low hundreds of megabytes, then round up so the count is a multiple of total executor cores, giving every core work and leaving room for a second wave. Too few partitions means huge tasks that spill or die; too many means thousands of tiny blocks, per-task scheduling overhead, and a small-file problem if you write straight out. Note the RDD API ignores this key and uses `spark.default.parallelism`. Adaptive execution can coalesce over-provisioned partitions after the fact, so erring slightly high is safer than erring low.
code
text · 4 lines== Physical Plan ==
*(5) HashAggregate(keys=[user_id#12], functions=[sum(amount#15)])
+- Exchange hashpartitioning(user_id#12, 200), ENSURE_REQUIREMENTS, [plan_id=88]
+- *(4) HashAggregate(keys=[user_id#12], functions=[partial_sum(amount#15)])go deeper
Know that the setting controls how many partitions a DataFrame shuffle produces, that the default is 200, and that big jobs usually need far more.
Derive a number: shuffled bytes divided by a target partition size, rounded to a multiple of total cores. Name the failure modes at both extremes and note that the RDD API uses a different key.
Set it per job from measured stage statistics, recognise spill and long tasks as the too-low signature and fetch overhead plus small files as the too-high one, and explain how adaptive coalescing changes the safe direction to err in.
Decide what the platform default should be and how teams override it, weighing under-sized jobs failing loudly against over-sized ones wasting scheduler capacity and producing small files that every downstream reader pays for.
## What the setting controls `spark.sql.shuffle.partitions` is the number of output partitions of every exchange in a SQL or DataFrame plan: joins, aggregations, window operations, explicit repartitions by expression. You can see it baked into the physical plan as the second argument of the `Exchange hashpartitioning(...)` node. Its default is **200**, and that number has nothing to do with your data — the same 200 applies to a 10 MB join and a 10 TB one. It is the most consequential single default in Spark, which is why the question is asked so often. One trap worth stating up front: the RDD API does not read this key at all. `reduceByKey` and friends fall back to `spark.default.parallelism` (or an explicitly passed partition count). Mixing the two APIs and tuning only one is a common source of "my setting did nothing". ## Sizing from data, not from habit The usable rule has two constraints that must both hold. **Constraint 1 — bytes per partition.** Take the stage's shuffle-write size from a previous run (the stage detail page reports it) and divide by a target per-partition size. A few hundred megabytes per post-shuffle partition is a workable target band: large enough that per-task overhead is negligible, small enough to sort within a task's slice of execution memory. Spark's own adaptive machinery uses a target on this order for the partitions it coalesces, which is a reasonable anchor. If a stage shuffles 2 TB and you want roughly 256 MB partitions, you need on the order of 8,000 partitions — 40× the default. **Constraint 2 — cores.** The count should be at least the total number of executor cores, or the cluster idles; making it a small multiple (two to four waves) helps the scheduler smooth over uneven task durations, since a straggler at the end of a single wave leaves everything else waiting. When the two constraints disagree — a small dataset on a large cluster — parallelism wins for latency but do not chase it past the point where tasks do less work than they cost to schedule. ## What goes wrong when it is too low Each task must sort and hold a huge partition. Symptoms, in order of appearance: heavy spill in the stage metrics, then task durations in the tens of minutes, then executors killed for exceeding memory, then fetch failures as those executors take their shuffle files with them. Leaving 200 on a multi-terabyte join is the single most common cause of the "my Spark job dies on big data but works on the sample" report. ## What goes wrong when it is too high The cost is not symmetric but it is real. Every map task writes a block for every reduce partition, so blocks in the shuffle scale as map tasks × partitions; at very high counts the reduce side spends its time on per-block fetch overhead rather than on bytes, and fetch wait time balloons while bytes read stay modest. The driver's map-output bookkeeping grows with the same product. And if the stage's output is written directly to storage, you get one small file per partition — the small-file problem, paid for later by every reader. ## Where adaptive execution fits Adaptive query execution, enabled by default since Spark 3.2 via `spark.sql.adaptive.enabled`, measures actual shuffle-write statistics at runtime and can coalesce an over-provisioned partition count down to a sensible one before the reduce stage runs. That changes the risk calculus: setting the number too high is now partially self-correcting, while setting it too low is not — nothing can split a partition you never created (skew handling aside). So with adaptive execution on, the practical advice is to set the count generously and let coalescing trim it, rather than to hand-tune a single perfect number. ## Per-job rather than global One cluster-wide value cannot suit every job, and even inside one job different stages shuffle wildly different volumes. Set it per job from that job's own data volume — `spark.conf.set` at runtime works between actions — and treat any platform-wide default as a starting point for the median job, not a governance decision. ## What interviewers are checking That you know the default is 200 and why that is arbitrary; that you can produce a sizing method involving both data volume and core count; that you can name the failure mode at each end; and that you do not answer "just set it to 2000" without saying what the number came from.
- Why is setting the partition count too high less dangerous than setting it too low?Too low produces oversized partitions that spill, run for tens of minutes and get executors killed, and nothing at runtime can split them. Too high produces many small tasks and blocks, which costs overhead — but adaptive execution measures the real shuffle sizes and coalesces excess partitions before the reduce stage. The self-correction runs in one direction only.
- Does spark.sql.shuffle.partitions affect RDD operations like reduceByKey?No. It applies only to SQL and DataFrame exchanges. RDD shuffles use the partition count passed to the operation, or `spark.default.parallelism` when none is given. Jobs that mix both APIs need both settings considered; tuning only the SQL key and seeing no change on an RDD stage is a frequent confusion.
- How would you pick the number without a prior run to measure?Estimate the shuffled volume from the input size and the plan — a join or aggregation typically shuffles on the order of the projected columns you keep — divide by a target of a few hundred megabytes, and round up to a multiple of total cores. Then run once, read the stage's actual shuffle-write size and spill metrics, and correct.
saying these in an interview costs you the question
- Says 200 is a tuned default suited to most jobs
- Sets one global value for every job on the platform
- Believes raising it always improves performance
- Thinks it controls RDD shuffles as well
- Ignores output file count when the stage writes to storage