CertSafari
    Databricks Certified Associate Developer for Apache Spark· Lessons

    Domain 4 · Lesson 22/32

    Spark data skew and shuffle tuning with Adaptive Query Execution

    Implement performance tuning strategies & optimize cluster utilization, including partitioning, repartitioning, coalescing, identifying data skew, and reducing shuffling

    9 min read
    3.12% of exam
    4 sources
    Published 3 Oct 2026
    Docs as of 30 Sep 2026

    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. 1.In Summary Metrics, compare the Max duration with the 75th percentile duration
    2. 2.If there is neither spill nor skew, check whether the stage is I/O bound
    3. 3.Check the stage details at the top of the page for spill

    Sources12

    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.

    Databricks AQE coalescing settings and their defaults
    PropertyDefaultEffect
    spark.sql.adaptive.coalescePartitions.enabledtrueTurns partition coalescing on or off
    spark.sql.adaptive.advisoryPartitionSizeInBytes64MBTarget size; coalesced partitions get close to it but no bigger
    spark.sql.adaptive.coalescePartitions.minPartitionSize1MBCoalesced partitions are no smaller than this
    spark.sql.adaptive.coalescePartitions.minPartitionNum2x no. of cluster coresNot 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?

    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?

    Sources34

    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.

    Databricks AQE skew join settings
    PropertyDefaultRole
    spark.sql.adaptive.skewJoin.enabledtrueTurns skew join handling on or off
    spark.sql.adaptive.skewJoin.skewedPartitionFactor5Multiplied by the median partition size
    spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes256MBAbsolute 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?

    Sources34

    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.

    Checkpoint 6 of 6· Exam question

    Which statement correctly distinguishes `repartition()` from `coalesce()` on a DataFrame?

    Sources34

    Exam traps

    Each one states something that sounds right. Open it to see what is actually true.

    1. 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. 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. 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. 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. 3.
      “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
    3. 4.
      “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

    Spotted a mistake, or was something unclear? Tell us.