What you will be able to do
- Identify skew in a long-running stage from the Spark UI Summary Metrics
- Explain how AQE coalesces small post-shuffle partitions, and which settings control the result
- Apply AQE's two-part test for a skewed partition
- Reduce shuffle cost with broadcast joins, and check AQE's changes in the query plan
1.Identifying skew in a long-running stage
In data skew, one or a few tasks run much longer than the rest. While those few tasks finish, most cores sit idle, so skew causes both slow jobs and poor cluster utilization. The Databricks Spark UI guide gives an order for checking a slow stage. Open the stage page and look at the details at the top for spill. Spill happens when Spark runs low on execution memory and starts writing data to disk, and it is most common during shuffles.
If a stage has neither spill nor skew, the guide sends you to check whether the stage is I/O bound. Skew can also cause failures, not just slowness. The Databricks memory guide lists skew among the causes of memory problems, along with too few shuffle partitions. Both come back to how data is spread across partitions.
Checkpoint 1 of 6· Put it in order
Put the Databricks guide's checks for a long-running stage in order.
- 1.In Summary Metrics, compare the Max duration with the 75th percentile duration
- 2.If there is neither spill nor skew, check whether the stage is I/O bound
- 3.Check the stage details at the top of the page for spill
The guide says to look for spill first, then skew in Summary Metrics. Only when a stage shows neither does it move on to high I/O.
“The first thing to look for in a long-running stage is whether there's spill.”Source: docs.databricks.com
2.AQE coalescing of post-shuffle partitions
A fixed spark.sql.shuffle.partitions can't suit every query. Adaptive Query Execution (AQE) helps by re-optimizing the query while it runs. At the end of each shuffle or broadcast exchange, Spark has accurate statistics for the data it just produced, and it uses them to adjust the rest of the plan. AQE is on by default (spark.databricks.optimizer.adaptive.enabled on Databricks, spark.sql.adaptive.enabled in Apache Spark since 3.2.0). It applies to non-streaming queries that contain at least one exchange or subquery. Joins, aggregates and window functions usually create the exchange.
One AQE feature merges small post-shuffle partitions into reasonably sized ones. This matters because very small tasks have worse I/O throughput and spend proportionally more time on scheduling and task setup. In Apache Spark, the practical approach is to set a large enough initial partition count with spark.sql.adaptive.coalescePartitions.initialPartitionNum and let AQE merge partitions down at runtime. This feature only combines partitions. It never makes them smaller. Splitting oversized partitions is done by a separate skew feature, covered in the next section.
| Property | Default | Effect |
|---|---|---|
| spark.sql.adaptive.coalescePartitions.enabled | true | Turns partition coalescing on or off |
| spark.sql.adaptive.advisoryPartitionSizeInBytes | 64MB | Target size; coalesced partitions get close to it but no bigger |
| spark.sql.adaptive.coalescePartitions.minPartitionSize | 1MB | Coalesced partitions are no smaller than this |
| spark.sql.adaptive.coalescePartitions.minPartitionNum | 2x no. of cluster cores | Not recommended; setting it overrides minPartitionSize |
Open-source Spark adds one setting that affects utilization: spark.sql.adaptive.coalescePartitions.parallelismFirst, which defaults to true. While it is true, Spark ignores the 64 MB target. It respects only the 1 MB minimum, to keep as much parallelism as possible. On a busy cluster, the docs recommend setting it to false so that the cluster isn't flooded with small tasks.
Checkpoint 2 of 6· Check yourself
A shared, busy Apache Spark cluster runs many queries that end up with lots of tiny post-shuffle tasks, even though AQE is on. Which change does the tuning guide recommend?
While parallelismFirst is true, Spark ignores the 64 MB advisory size and only respects the 1 MB minimum. Setting it to false on a busy cluster makes AQE merge partitions toward the target size.
“It's recommended to set this config to false on a busy cluster to make resource utilization more efficient (not many small tasks).”Source: spark.apache.org
Checkpoint 3 of 6· Exam question
In the Spark UI, a job's shuffle stage shows 199 tasks finishing in under 10 seconds while a single task runs for 25 minutes and dominates the stage's wall-clock time. The stage performs a `groupBy("customer_id").agg(sum("amount"))` where one `customer_id` accounts for a large share of all rows. What is this symptom, and which change directly addresses its root cause?
Correct answer: A — This is data skew from an overrepresented key; salt the skewed key with a random suffix, aggregate on the salted key, then aggregate the partial results again to combine them.
- A. One massively overrepresented key sending nearly all rows to a single reduce task is the textbook definition of data skew, and salting spreads that key's rows across many synthetic sub-keys before a second aggregation recombines them, which is the standard fix that attacks the uneven data distribution directly.
- B. More executor memory can help a task that is failing with out-of-memory errors, but it does not rebalance how many rows land on each task, so the one oversized partition would still take far longer than the others.
- C. Compacting small files improves read throughput and file listing overhead, but it has no effect on how shuffle output is partitioned by aggregation key, so the skewed key would still dominate a single task.
- D. There is no broadcast exchange in this plan since `groupBy` with an aggregation triggers a shuffle, not a broadcast, so raising the broadcast timeout has no bearing on this stage's runtime at all.
- E. Trimming unused columns can reduce shuffle payload size somewhat, but it does not change which rows are routed to which reduce task, so the key causing the imbalance would still overwhelm one task.
3.AQE skew join handling
Skew in a join produces the straggler tasks you saw in Summary Metrics. AQE handles skew in sort merge joins and shuffle hash joins. It splits each skewed task into roughly evenly sized tasks, replicating the matching data from the other side when needed. Before this feature existed, skew handling needed hints. A partition counts as skewed only if it passes two tests at once.
| Property | Default | Role |
|---|---|---|
| spark.sql.adaptive.skewJoin.enabled | true | Turns skew join handling on or off |
| spark.sql.adaptive.skewJoin.skewedPartitionFactor | 5 | Multiplied by the median partition size |
| spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes | 256MB | Absolute size the partition must also exceed |
Checkpoint 4 of 6· Check yourself
The median shuffle partition in a join is 100 MB and the largest is 400 MB. All skew settings are at their defaults. Does AQE treat the 400 MB partition as skewed?
Both conditions must hold. 400 MB passes the 256 MB threshold but fails the 5 × median test, because 5 × 100 MB is 500 MB.
“A partition is considered skewed when both (partition size > skewedPartitionFactor * median partition size) and (partition size > skewedPartitionThresholdInBytes) are true.”Source: docs.databricks.com
The advisory size (64 MB) is the size AQE aims for when it merges or splits partitions. The 256 MB threshold should be well above it, so that AQE splits only partitions that are clearly oversized, not ordinary ones near the target.
4.Reducing shuffles with broadcast joins, and checking the plan
The cheapest shuffle is one that never happens. A sort merge join shuffles and sorts both sides. A broadcast hash join sends the small table to every worker instead, so the large side doesn't move. The planner broadcasts a table when its estimated size is below spark.sql.autoBroadcastJoinThreshold, which defaults to 10 MB. Setting it to -1 disables broadcasting.
Size estimates made before the query runs are often wrong, so AQE checks again at runtime. If the runtime statistics show one side of a sort merge join is small enough, AQE switches the join to a broadcast hash join. On Databricks, the runtime threshold is spark.databricks.adaptive.autoBroadcastJoinThreshold, which defaults to 30MB. The switch is less efficient than planning a broadcast from the start, because the shuffle has already been written. It still avoids sorting both sides. With spark.sql.adaptive.localShuffleReader.enabled, Spark also reads the shuffle files locally instead of over the network.
To see what AQE did, look at the plan. Queries where AQE applies have AdaptiveSparkPlan nodes. SQL EXPLAIN doesn't run the query, so it always shows the initial plan. Each AQE change only appears in the current or final plan, as a specific node.
Checkpoint 5 of 6· Match them up
Match each AQE optimization to the evidence it leaves in the current or final plan.
Tap a term, then the definition that fits it.
Each AQE feature leaves its own marker in the plan. You can only see these markers by comparing the current or final plan with the initial plan.
“Dynamically handle skew join: node SortMergeJoin with field isSkew as true.”Source: docs.databricks.com
Checkpoint 6 of 6· Exam question
Which statement correctly distinguishes `repartition()` from `coalesce()` on a DataFrame?
Correct answer: A — `repartition()` always triggers a full shuffle and can increase or decrease the number of partitions, while `coalesce()` avoids a full shuffle and can only decrease the partition count.
- A. This matches Spark's actual behavior: `repartition()` redistributes all rows across the cluster via a full shuffle and supports moving to either more or fewer partitions, while `coalesce()` merges existing partitions on the same executors to reduce the count without a full shuffle.
- B. This reverses the two methods' actual behavior; it is `coalesce()` that avoids a shuffle and only reduces partition count, while `repartition()` is the one that shuffles and can move in either direction.
- C. `coalesce()` specifically avoids a full shuffle in the common case of reducing partitions, and it does not guarantee balanced output sizes since it merges adjacent existing partitions rather than rehashing rows evenly.
- D. Both methods operate identically regardless of whether Spark runs in local mode or on a multi-node cluster; neither is restricted to a single deployment mode.
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.AQE's partition-coalescing feature also breaks large shuffle partitions into smaller ones.Why is that wrong?
Coalescing only combines small partitions into reasonably sized ones. Skewed partitions are split by the separate skew join feature, and only in sort merge and shuffle hash joins.
Covered in AQE coalescing of post-shuffle partitions
2.Any partition larger than 256 MB counts as skewed.Why is that wrong?
Spark treats a partition as skewed only if it exceeds both skewedPartitionThresholdInBytes and skewedPartitionFactor (default 5) times the median partition size.
Covered in AQE skew join handling
3.Running SQL EXPLAIN shows whether AQE turned a sort merge join into a broadcast join.Why is that wrong?
EXPLAIN doesn't execute the query, so it shows only the initial plan. AQE's changes appear in the current or final plan, after the query has run.
Covered in Reducing shuffles with broadcast joins, and checking the plan
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“Skew is when one or just a few tasks take much longer than the rest. This results in poor cluster utilization and longer jobs.”
↩︎ Identifying skew in a long-running stage“Spill is what happens when Spark runs low on execution memory.”
↩︎ Identifying skew in a long-running stage“If the stage doesn't have spill or skew, see Spark stage high I/O for the next steps.”
↩︎ Identifying skew in a long-running stage“If the Max duration is 50% more than the 75th percentile, you may be suffering from skew.”
↩︎ Prediction“The first thing to look for in a long-running stage is whether there's spill.”
↩︎ Checkpoint - 2.
“Too few shuffle partitions”
↩︎ Identifying skew in a long-running stage - 3.https://docs.databricks.com/aws/en/optimizations/aqeOfficial docs
“Very small tasks have worse I/O throughput and tend to suffer more from scheduling overhead and task setup overhead.”
↩︎ AQE coalescing of post-shuffle partitions“Contain at least one exchange (usually when there's a join, aggregate, or window), one sub-query, or both.”
↩︎ AQE coalescing of post-shuffle partitions“Dynamically handles skew in sort merge join and shuffle hash join by splitting (and replicating if needed) skewed tasks into roughly evenly sized tasks.”
↩︎ AQE skew join handling“Dynamically changes sort merge join into broadcast hash join.”
↩︎ Reducing shuffles with broadcast joins, and checking the plan“The threshold to trigger switching to broadcast join at runtime.”
↩︎ Reducing shuffles with broadcast joins, and checking the plan“AQE-applied queries contain one or more AdaptiveSparkPlan nodes, usually as the root node of each main query or sub-query.”
↩︎ Reducing shuffles with broadcast joins, and checking the plan“Dynamically coalesces partitions (combine small partitions into reasonably sized partitions) after shuffle exchange.”
↩︎ Exam trap 1“A partition is considered skewed when both (partition size > skewedPartitionFactor * median partition size) and (partition size > skewedPartitionThresholdInBytes) are true.”
↩︎ Exam trap 2“As SQL EXPLAIN does not execute the query, the current plan is always the same as the initial plan”
↩︎ Exam trap 3“A partition is considered skewed when both (partition size > skewedPartitionFactor * median partition size) and (partition size > skewedPartitionThresholdInBytes) are true.”
↩︎ Checkpoint“Dynamically handle skew join: node SortMergeJoin with field isSkew as true.”
↩︎ Checkpoint - 4.https://spark.apache.org/docs/latest/sql-performance-tuning.htmlSecondary source
“Spark can pick the proper shuffle partition number at runtime once you set a large enough initial number of shuffle partitions”
↩︎ AQE coalescing of post-shuffle partitions“Ideally, this config should be set larger than spark.sql.adaptive.advisoryPartitionSizeInBytes.”
↩︎ AQE skew join handling“spark.sql.autoBroadcastJoinThreshold | 10485760 (10 MB)”
↩︎ Reducing shuffles with broadcast joins, and checking the plan“we can avoid sorting both join sides and read shuffle files locally to save network traffic”
↩︎ Reducing shuffles with broadcast joins, and checking the plan“It's recommended to set this config to false on a busy cluster to make resource utilization more efficient (not many small tasks).”
↩︎ Checkpoint