skip to content

In Spark AQE, when does spark.sql.adaptive.skewJoin.enabled split a shuffle partition?

level: seniorimportance: should knowfreq 56%

answer

  1. it only looks after the shuffle has been written
  2. two conditions, not one
  3. five times what, exactly?
  4. the other side has to come along to each piece
  5. joins only — your groupBy is on its own

basics

~20 s

AQE splits a shuffle partition when it is both more than five times the median partition size and larger than 256 MB. It cuts the oversized partition into smaller pieces and replicates the matching partition from the other side so the pieces join in parallel.

solid answer

~40 s

After a join's shuffle map stages finish, Spark knows every partition's real size. A partition is treated as skewed when it exceeds `spark.sql.adaptive.skewJoin.skewedPartitionFactor` (5.0) times the **median** partition size **and** exceeds `spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes` (256 MB) — both conditions, so healthy stages with small partitions are left alone. Spark then splits that partition's map-output blocks into chunks of roughly `spark.sql.adaptive.advisoryPartitionSizeInBytes` and replicates the other side's corresponding partition to each chunk, so several tasks share the hot key's work instead of one. It needs `spark.sql.adaptive.enabled` and `spark.sql.adaptive.skewJoin.enabled` — both true by default since Spark 3.2 — and it targets sort-merge joins. It does nothing for a skewed `groupBy`.

code

properties · 5 lines
properties
spark.sql.adaptive.enabled=true
spark.sql.adaptive.skewJoin.enabled=true
spark.sql.adaptive.skewJoin.skewedPartitionFactor=5.0
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes=256m
spark.sql.adaptive.advisoryPartitionSizeInBytes=64m

go deeper

for a junior

Know that Spark 3 can notice an oversized shuffle partition in a join at runtime and break it into several tasks, and that this is on by default rather than something you must code.

for a middle

Explain both conditions a partition must meet, name the two configs and their defaults, and describe the split-and-replicate mechanic that keeps the join result correct.

for a senior

Diagnose why the rule did not fire on a real job — adaptive execution off, an aggregation rather than a join, a stage below the byte threshold — and decide between adjusting the thresholds and changing the data distribution.

for a principal

Decide how much skew handling belongs in engine configuration versus in data modelling. Reason about whether a recurring hot key should be designed out of the pipeline rather than papered over by a runtime rule at every consumer.

## Where the rule sits Adaptive Query Execution cuts the physical plan at shuffle boundaries. When a query stage completes, its map-output statistics give the exact byte size of every shuffle partition, and AQE re-optimises the rest of the plan with those numbers in hand. The skew-join rule is one of those re-optimisations. It runs after the join's inputs have been shuffled and before the join's reducer tasks are scheduled. Two switches gate it: `spark.sql.adaptive.enabled` (the umbrella, `true` by default since Spark 3.2) and `spark.sql.adaptive.skewJoin.enabled` (`true` by default since 3.0). Turning off the umbrella turns off everything below it — a common cause of "AQE didn't help" on clusters where someone disabled adaptive execution years ago and never revisited it. ## The two thresholds, and why there are two A partition is considered skewed only if **both** are true: 1. its size is greater than `spark.sql.adaptive.skewJoin.skewedPartitionFactor` × the median partition size — default factor `5.0`; 2. its size is greater than `spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes` — default `256MB`. The factor alone would fire constantly on healthy stages: if the median partition is 2 MB, a 12 MB partition is six times the median and completely harmless. The absolute floor stops Spark from splitting partitions that were never going to be slow. The trade-off is the mirror image: a stage whose partitions are all under 256 MB gets no skew handling even when one is twenty times its neighbours, which is a real reason a skewed job sometimes goes untouched. The documentation also advises keeping the threshold above `spark.sql.adaptive.advisoryPartitionSizeInBytes`, so the rule does not split into pieces at the same scale as the thing it is judging. ## Splitting and replicating The rule does **not** re-hash the key. A shuffle partition is physically a set of blocks, one per map task. Splitting means handing subsets of those blocks to different reducer tasks, targeting roughly the advisory partition size per piece. Since each piece now contains only some of the hot key's rows, the matching partition on the other side of the join must be given to *every* piece — that is the replication. For a hot key joined against a handful of dimension rows this is cheap; the replicated side is by definition the non-skewed one. The result is that a partition that would have been one 9 GB task becomes, say, thirty tasks of a few hundred megabytes, each joining its slice of the fact rows against the same replicated dimension rows. Output is identical, because every original pair of matching rows still meets exactly once: each fact row lives in exactly one piece, and the dimension rows are present in all of them. ## Where it applies, and where it does not - **Sort-merge joins** are the documented target. The rule needs post-shuffle partitions on both sides to measure and replicate. - **Broadcast hash joins** are untouched — there is no shuffle on the large side, and no concentration to fix. - **Aggregations are not covered.** A skewed `groupBy` sends every row of the hot group to one reducer task, and no splitting is possible without changing the aggregation itself, since the group's final value depends on all of its rows. That is what two-phase salted aggregation is for. - **Map-side skew is not covered.** If one input *file* or one Kafka partition is enormous, the imbalance exists before any shuffle statistic is produced. - The rule only fires when the plan shape lets it. `spark.sql.adaptive.forceOptimizeSkewedJoin` (default `false`) forces the optimisation even when it introduces an extra shuffle; leave it off unless you have measured that the extra exchange costs less than the straggler. ## Related shuffle-side setting AQE's decision depends on the map-output statistics being accurate. To keep `MapStatus` small, Spark compresses per-block sizes above a certain number of blocks, recording exact sizes only for blocks above `spark.shuffle.accurateBlockThreshold`. `spark.shuffle.accurateBlockSkewedFactor` (default `-1.0`, i.e. off) additionally records a block accurately when it exceeds that factor times the median block size; the documentation suggests setting it to the same value as `spark.sql.adaptive.skewJoin.skewedPartitionFactor`. On a job where you know skew exists and AQE seems blind to it, this is worth checking. ## What good tuning looks like Mostly, leave the defaults alone and verify from the Spark UI that the straggler disappeared. When it did not: check that adaptive execution is on at all; check whether the operator is a join or an aggregation; check whether the skewed partition actually clears 256 MB, and lower `skewedPartitionThresholdInBytes` if the whole stage runs at a smaller scale; and if the hot key is enormous — a large fraction of a multi-terabyte fact table — accept that splitting after the shuffle still means the shuffle wrote one giant partition, and change the data distribution with salting or a broadcast instead.

  • Why does the rule need both a factor and an absolute byte threshold?
    Either one alone misfires. The factor alone would flag partitions in a perfectly healthy stage — a 12 MB partition next to a 2 MB median is 6x the median and harmless. The byte floor alone would flag every large partition in a uniformly large stage, where nothing is actually imbalanced. Requiring both keeps splitting for partitions that are genuinely oversized and genuinely slow.
  • When would you set spark.sql.adaptive.forceOptimizeSkewedJoin to true?
    It defaults to false because forcing the optimisation can introduce an extra shuffle. Turn it on when a job repeatedly stalls on one straggler task in a join whose plan otherwise keeps the rule from firing, and you have measured that the additional exchange costs less than the straggler it removes. Treat it as a per-job setting, not a cluster default.
  • Does the skew rule fire on a join that was already planned as a broadcast hash join?
    No. A broadcast hash join has no shuffle on the large side, so there are no post-shuffle partitions to measure or split — the hot key stays inside whichever scan partition it was read into and joins locally. The rule exists for shuffle-based joins, which is precisely where hashing a key concentrates it into one task.

saying these in an interview costs you the question

  • Thinks AQE rebalances skew in any operator, including aggregations
  • Says a partition is skewed purely for exceeding 256 MB
  • Forgets the other side's partition is replicated to each split
  • Believes skew splitting works with adaptive execution disabled
  • Assumes AQE detects skew before the shuffle is written

context