Skip to content
Great engineers know 3 min read · skewJoin.enabled 1 practice problem ↓

AQE skew-join handling

How Spark detects oversized partitions in sort-merge joins and splits them automatically.

You will learn

  • How AQE decides that a join partition is skewed
  • How it splits a skewed partition without breaking the join
  • Which join types it supports
  • How to tune and verify it

Read first

Comfortable with these? Read on.

TL;DR After the shuffle of a sort-merge join, AQE marks a partition as skewed when it is more than 5 times the median size and larger than 256 MB. It splits that partition into smaller pieces and duplicates the matching partition of the other side, so each piece joins independently.

When a partition counts as skewed

Both conditions must hold:

ConfigDefaultCondition
spark.sql.adaptive.skewJoin.skewedPartitionFactor5Partition size is more than 5 × the median partition size
spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes256MBAnd larger than 256 MB
spark.sql.adaptive.skewJoin.enabledtrueThe feature is on (needs AQE on)

The threshold stops AQE from splitting partitions that are relatively large but absolutely small: a 20 MB partition next to 1 MB ones is not worth the overhead.

How the split works

Before · one huge task

Orders joined to customers by customer_id. After the shuffle, partition 17 of orders is 12 GB (one huge customer); the median is 150 MB. Customers' partition 17 is 40 MB. Without AQE, one task joins 12 GB.

Split · by map output

AQE splits orders partition 17 into about 12 GB / advisory size pieces, around 190 pieces at 64 MB. The split follows the boundaries of the map outputs that wrote the shuffle data, so it needs no extra shuffle.

Replicate · the other side

Each piece still needs every customer row of partition 17 to join correctly, so customers partition 17 (40 MB) is read by each of the 190 tasks. That duplication is the cost.

After · 190 tasks

The 12 GB join runs as 190 parallel tasks. The result is identical; the stage time drops from the time of the slowest task to roughly the time of an average one.

Supported joins

Join typeWhich side can be split
InnerEither or both
Left outer, left semi, left antiLeft side only
Right outerRight side only
Full outerNot supported

The rule mirrors broadcasting: the side whose unmatched rows must be preserved can be split, because each piece is joined against the full other side. Splitting the non-preserved side would make unmatched rows ambiguous.

Tuning and verifying

  • If a skewed partition is not split, check both conditions: on a cluster where every partition is large, 5 × median may never be reached; lower the factor.
  • Some plans, such as a join whose output is immediately re-shuffled the same way, cannot be split without an extra shuffle. spark.sql.adaptive.forceOptimizeSkewedJoin (3.3+) lets AQE add that shuffle.
  • In the final plan (SQL tab of the Spark UI), the join shows skew=true and an AQEShuffleRead node reporting the number of skewed partitions and splits.
PySparkSpark SQL
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionFactor", "3")
spark.conf.set("spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes", "128MB")
SET spark.sql.adaptive.skewJoin.skewedPartitionFactor = 3;
SET spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes = 128MB;
Note: a single key cannot be split below one map task's output for that key, and replication of the other side costs I/O. For extreme skew, salting or handling the hot key separately can still beat AQE.

Common mistakes

Expecting it to fix skewed aggregations

It only handles joins.

Assuming full outer joins are covered

They are not.

Leaving AQE off

Skew-join handling needs spark.sql.adaptive.enabled.

Key takeaways

  • A partition is skewed when it exceeds 5 × the median and 256 MB (defaults).
  • AQE splits it and replicates the other side's matching partition.
  • Only sort-merge-style joins; outer joins only on the preserved side; never full outer.
  • Verify with skew=true in the final plan.

Check yourself

3 questions

1. With defaults, median partition 100 MB: is a 400 MB partition split?

Show the answer

No: it is not more than 5× the median. 400 MB is 4× the median, below the factor of 5.

2. What does AQE do to the non-skewed side's matching partition?

Show the answer

Reads it once per split of the skewed side. Each piece needs all matching rows, so that partition is read by every split task.

3. In a left outer join, which side can AQE split?

Show the answer

Left. The preserved (left) side can be split; each piece joins against the full right partition.

Practice it

Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.

Solve: Find Hot Keys and Plan Salting →

Go deeper

Primary sources: AQE: optimizing skew join