skip to content

What does spark.sql.adaptive.coalescePartitions.enabled do to a stage's task count?

level: middleimportance: should knowfreq 68%

answer

  1. the 200 default is set before the data is seen
  2. real map-output sizes arrive after the shuffle writes
  3. neighbouring partitions get merged, not re-hashed
  4. it can shrink the count but never grow it
  5. one default quietly ignores the 64 MB target

basics

~20 s

With Adaptive Query Execution on, Spark reads the finished map stage's real output sizes and merges contiguous small shuffle partitions into fewer, larger reducer tasks. It can only reduce the post-shuffle partition count, never raise it.

solid answer

~40 s

`spark.sql.shuffle.partitions` is a static number — 200 by default — so a stage that produces only a few hundred megabytes still gets 200 tiny tasks, and one that produces terabytes gets 200 huge ones. When AQE is on (the default since Spark 3.2) and `spark.sql.adaptive.coalescePartitions.enabled` is true, Spark waits for the map stage's real output statistics and merges **contiguous** partition ranges into a single reducer task. The target is `spark.sql.adaptive.advisoryPartitionSizeInBytes`, 64 MB by default — but `spark.sql.adaptive.coalescePartitions.parallelismFirst` also defaults to true, which ignores that advisory size and enforces only `spark.sql.adaptive.coalescePartitions.minPartitionSize` (1 MB), to protect parallelism. Set it to false on a busy cluster if you actually want 64 MB tasks. Coalescing never splits, so set a deliberately high `spark.sql.adaptive.coalescePartitions.initialPartitionNum` and let AQE come down.

code

properties · 6 lines
properties
spark.sql.adaptive.enabled=true
spark.sql.adaptive.coalescePartitions.enabled=true
spark.sql.adaptive.coalescePartitions.initialPartitionNum=2000
spark.sql.adaptive.advisoryPartitionSizeInBytes=64m
# without this, the 64m target is ignored in favour of parallelism
spark.sql.adaptive.coalescePartitions.parallelismFirst=false

go deeper

for a junior

Know that spark.sql.shuffle.partitions defaults to 200 and that AQE, on by default since Spark 3.2, can reduce that number at runtime so small stages do not run hundreds of near-empty tasks.

for a middle

Explain where the runtime statistics come from, why only contiguous partitions are merged, and why coalescing can shrink but never grow the partition count. Name the advisory size and the parallelismFirst default.

for a senior

Show the tuning pattern you actually apply: a high initialPartitionNum, parallelismFirst turned off on a shared cluster, and the check in the Spark UI that confirms task sizes landed where you wanted them.

for a principal

Own the cluster-wide defaults. Decide whether parallelism or resource efficiency wins on a shared cluster, and set the AQE baseline centrally so hundreds of jobs stop shipping hand-tuned partition counts that rot.

## The problem it solves Every shuffle in Spark SQL produces a fixed number of partitions, taken from `spark.sql.shuffle.partitions`, whose default is 200. That number is chosen before the query runs and applies to every shuffle in it. It is almost always wrong for at least one stage: a filtered branch that emits 300 MB gets 200 partitions of 1.5 MB each — 200 tasks whose scheduling and shuffle-fetch overhead dwarfs their work — while a stage that emits 2 TB gets 200 partitions of 10 GB each, which spill and often fail. Before Spark 3, tuning this by hand per job (and sometimes accepting a compromise between the stages inside one job) was a routine part of Spark work. ## What coalescing does Adaptive Query Execution breaks the physical plan at shuffle boundaries into *query stages*. When a map stage finishes, its `MapStatus` output carries the exact size of every shuffle partition — no estimation involved. The coalescing rule reads those sizes and, instead of scheduling one reducer task per shuffle partition, assigns each reducer task a **contiguous range** of partitions whose combined size approaches the target. Two hundred 1.5 MB partitions might be served by five reducer tasks reading 40 partitions each. Contiguity matters: a reducer task fetches a range of partition ids from every map output, so merging neighbours is a pure bookkeeping change on the read side. No data is re-hashed, no second shuffle is written, and the rows still land in a task that owns whole partitions — so grouping and join semantics are untouched. ## The configuration surface - `spark.sql.adaptive.enabled` — the umbrella switch, `true` by default since Spark 3.2. Nothing below works without it. - `spark.sql.adaptive.coalescePartitions.enabled` — `true` by default. - `spark.sql.adaptive.advisoryPartitionSizeInBytes` — the target size, 64 MB by default. It is also the target used when splitting skewed partitions. - `spark.sql.adaptive.coalescePartitions.parallelismFirst` — `true` by default. When true, Spark **ignores** the advisory size and only enforces the minimum size, so it merges just enough to eliminate trivially small partitions and keeps parallelism high. The documentation recommends setting it to `false` on a busy cluster, where many small tasks waste more than they gain. - `spark.sql.adaptive.coalescePartitions.minPartitionSize` — 1 MB by default; the floor honoured when the advisory size is being ignored. - `spark.sql.adaptive.coalescePartitions.initialPartitionNum` — the partition count the shuffle *starts* with; falls back to `spark.sql.shuffle.partitions` when unset. That `parallelismFirst` default surprises people: they set the advisory size to 128 MB, see tasks of a few megabytes, and conclude AQE is broken. It is doing exactly what it was configured to do. ## Why it only merges The map side has already written its output into N buckets by the time AQE looks. Splitting one of those buckets into finer key groups would require re-partitioning the data — another shuffle. Merging neighbours requires nothing. That asymmetry drives the recommended tuning pattern: pick a generous `initialPartitionNum` (thousands, for a large job), let the shuffle write fine-grained buckets, and let coalescing produce sensible reducer tasks per stage. You trade a little map-side bookkeeping for the ability to size every stage correctly, instead of one number for all of them. (The exception is skew, where AQE *can* split a partition — but only for joins, and by splitting the map-output blocks of one oversized partition, not by re-hashing keys.) ## What it does not fix Coalescing balances *small* partitions; it does nothing about one partition being enormous. If a hot key produces a 9 GB bucket, coalescing leaves it alone (merging it with a neighbour would only make it worse) and the straggler task remains. The skew-join rule, salting, or broadcasting is what addresses that. It also is not the DataFrame `coalesce(n)` transformation. That is a user-invoked narrow transformation that reduces partitions in the *current* stage without a shuffle, and, because it collapses the stage's parallelism upstream, it can starve the computation feeding it. AQE coalescing is an automatic runtime decision about *post-shuffle* reducer tasks. Sharing a name is unfortunate; interviewers ask about it precisely because the two are confused. ## Seeing it work In the SQL tab, an AQE plan is rooted at an `AdaptiveSparkPlan` node whose `isFinalPlan` flag flips to `true` once execution has settled, and the shuffle-read nodes report how many partitions were coalesced. The simpler signal is the stage's task count: fewer tasks than `spark.sql.shuffle.partitions`, with sensible per-task shuffle read sizes, means coalescing fired.

  • Does coalescing remove the need to set spark.sql.shuffle.partitions at all?
    It removes the need to get it exactly right, but the value still sets the ceiling: AQE can only merge the partitions the shuffle actually produced. If 200 is far too few for a multi-terabyte stage, every merged task is still enormous. Set spark.sql.adaptive.coalescePartitions.initialPartitionNum high enough that the shuffle starts fine-grained, then let coalescing bring the task count down per stage.
  • Why does the REBALANCE hint exist when coalescing already runs?
    REBALANCE targets the query's output partitions rather than an intermediate stage: it asks Spark to make each output partition a reasonable size, splitting skewed ones and merging small ones, so a write does not emit thousands of tiny files or one giant one. It is best-effort, can take column names, and is ignored entirely when AQE is off.
  • How is this different from calling coalesce(10) on a DataFrame?
    DataFrame coalesce is a narrow transformation you request explicitly: it merges partitions in the current stage with no shuffle, and it reduces the parallelism of everything upstream in that stage, which can badly slow a heavy computation. AQE coalescing is an automatic, post-shuffle decision about reducer task counts and does not constrain the map side at all.

saying these in an interview costs you the question

  • Thinks AQE can raise the shuffle partition count at runtime
  • Believes 200 partitions is chosen per job rather than a static default
  • Assumes the 64 MB advisory size applies with parallelismFirst on
  • Says coalescing needs an extra shuffle to regroup the rows
  • Confuses AQE coalescing with the DataFrame coalesce transformation

context