skip to content

Skew and Adaptive Execution

Skew is when one key holds a disproportionate share of the data and one task holds up the whole stage. Adaptive Query Execution automates part of the fix at runtime, so interviewers ask both what you would do manually and what Spark 3 now does for you.

on this pageshow

explore

questions

7

In a Spark stage, which task metrics separate data skew from a slow executor?

level: middleimportance: must knowfreq 70%

answer

  1. compare tasks, not the whole stage
  2. look at the max-versus-median row
  3. is the straggler reading more, or just slower?
  4. one host, or one partition?
  5. group the key and count the rows

basics

~20 s

Compare the stage's per-task quantiles in the Spark UI. Skew shows one task reading many times the median task's shuffle bytes or records; a slow executor shows normal input sizes with long durations concentrated on one host.

solid answer

~40 s

Open the stage page in the Spark UI and read the summary metrics table, which reports min / 25th / median / 75th / max for task duration, shuffle read size and records, and spill. Skew looks like a max task whose **input** is 10-50x the median: the straggler is doing more work, not the same work more slowly, and it usually spills to disk while its peers spill nothing. A slow node looks different — input sizes are uniform across tasks, but the long tasks all land on one executor, so sort the task table by host and check GC time. If the stage is a join or aggregation, confirm at the data level with `df.groupBy(key).count().orderBy(desc("count"))`: one key holding a large share of the rows is the cause.

code

text · 5 lines
text
Summary Metrics for 200 Completed Tasks
Metric              Min     25th    Median  75th    Max
Duration            0.4 s   0.5 s   0.6 s   0.9 s   14 min
Shuffle Read Size   9 MB    11 MB   12 MB   13 MB   9.4 GB
Spill (Disk)        0.0 B   0.0 B   0.0 B   0.0 B   6.2 GB

go deeper

for a junior

Know that a Spark stage runs one task per partition and finishes only when its slowest task does. Be able to open the Spark UI's Stages tab and find how long the longest task took.

for a middle

Explain the summary metrics quantiles and why shuffle read size, not duration, is the metric that identifies skew. Distinguish skew, an unhealthy executor, and an under-partitioned stage from the same table.

for a senior

Show the full diagnosis loop on a real job: UI metrics, then a count of the key, then a remedy chosen for the operator involved. Be ready to explain why speculation and extra executors do not help.

for a principal

Own the detection side as well as the fix: event logging retention, a job-level check on max-versus-median task duration, and a policy for which pipelines get skew handling designed in rather than patched after an incident.

## What skew is in Spark terms A Spark stage is executed as one task per partition. Every wide transformation — a join, a `groupBy`, a `distinct`, a window — ends the previous stage with a shuffle: rows are written out bucketed by a hash of the grouping or join key, and each reducer task in the next stage fetches one bucket. If the key distribution is uniform, every bucket is about the same size and every task finishes at about the same time. If one key value holds a large share of the rows, all of those rows hash to the same bucket, and one task must read and process them alone. The stage cannot complete until that task does, so the stage's wall-clock time becomes that one task's time no matter how many executors are idle behind it. ## Where to look The Spark UI's **Stages** tab, opened on the slow stage, has a *Summary Metrics for Completed Tasks* table with quantile columns. The rows that matter are: - **Duration** — how long tasks took. - **Shuffle Read Size / Records** — how much each task actually pulled. - **Spill (Memory)** and **Spill (Disk)** — whether a task exceeded its execution memory and wrote sorted runs to local disk. - **Input Size / Records** for a scan stage rather than a shuffle stage. Below the summary is the full task list, sortable by duration and showing the executor and host for each task, plus the **Event Timeline**, which draws the tasks as bars and makes a single long bar in an otherwise dense block unmistakable. ## The diagnostic: input, not time Duration alone does not distinguish causes; the shape of the *input* does. - **Skew** — max shuffle read is many times the median while the median stays small. The task is slow because it holds more rows. Spill on that task and only that task is corroborating evidence: the partition no longer fits in execution memory. - **A slow or unhealthy executor** — inputs are even across tasks, but the tasks assigned to one executor run long. Look at the Executors tab for that host's GC time, its failed-task count, and whether other stages also drag on it. - **Too few partitions** — every task is big and slow together. The max/median ratio stays close to 1 and the whole quantile row shifts up. That is a sizing problem, not skew, and raising the partition count fixes it — whereas raising the partition count barely helps skew, because the hot key still hashes to exactly one partition no matter how many there are. ## Confirming the hot key Once the stage metrics point at skew, prove it in the data rather than guessing. Group the join or aggregation key, count, and order descending. Real-world offenders are recognisable: `NULL` or a placeholder such as `-1` / `"unknown"` standing in for a missing foreign key, a default tenant or system account, a bot user, or a date column where a backfill dumped every historical row onto one day. Nulls deserve special attention: in an inner equi-join they never match anything, so all that shuffled data produces no output rows at all. ## Two traps **Speculative execution does not fix skew.** `spark.speculation` re-launches a task that is running much slower than its peers on another executor and takes whichever copy finishes first. That is the right medicine for a straggler caused by a bad disk or a noisy neighbour, because the duplicate reads the same small input on healthy hardware. A skewed task's duplicate reads the same enormous partition and is just as slow — you have spent double the resources for nothing. **Adding executors does not fix skew either.** The stage's critical path is one task; extra cores sit idle waiting for it. This is the single most common wrong answer in an interview, and the reason the diagnosis step matters before any remedy. ## After the diagnosis What you do next depends on the operator. For a shuffle join, Adaptive Query Execution may already split the oversized partition for you; for a `groupBy` it will not, and you salt the key or pre-aggregate. If one side of the join is small, broadcasting it removes the shuffle and the skew together. If the hot key is a null or a sentinel, handling it separately upstream is usually cheaper than any engine-level trick. Keep `spark.eventLog.enabled` on in production so the same stage page is available in the history server after the run; skew is often noticed the morning after, when the live UI is gone.

  • 199 of the 200 tasks finish in a minute and one runs for an hour. Would speculative execution help?
    No. Speculation re-launches an unusually slow task elsewhere and keeps whichever copy finishes first, which helps when the cause is a bad disk or a noisy neighbour. A skewed task is slow because it holds far more data, so the speculative copy processes the same rows and takes just as long — you pay twice the resources for the same wall clock.
  • How would you tell skew apart from simply having too few shuffle partitions?
    Too few partitions makes every task big, so the whole quantile row shifts up and the max/median ratio stays near 1. Skew leaves the median task small and stretches only the tail. Raising the partition count fixes the first and barely moves the second, because the hot key still hashes to a single partition however many there are.
  • Where do you find per-task metrics for a job that already finished?
    The Spark history server renders the same stage page from the event log, provided spark.eventLog.enabled was true for that run and the log directory is readable. Its summary metrics table, task list and event timeline are identical to the live UI, which is why event logging is worth leaving on for scheduled pipelines.

A checkout line runs long either because one shopper has three full carts or because the cashier is new. You tell the two apart by looking at the carts, not at the clock.

saying these in an interview costs you the question

  • Reports the stage is slow without comparing per-task metrics
  • Blames the cluster before checking shuffle read size per task
  • Thinks adding executors speeds up a single overloaded task
  • Uses average task duration instead of max-versus-median spread
  • Expects speculative execution to rescue a skewed task

context

open as a page

In a Spark join, how does salting a hot key spread its rows across many tasks?

level: seniorimportance: must knowfreq 58%

basics

~20 s

You append a random 0..N-1 salt to the join key on the big side and replicate every row of the small side N times, once per salt value. The composite key now hashes into N partitions instead of one.

open as a page

In a Spark join, why does broadcasting the small side remove the shuffle and the skew?

level: middleimportance: should knowfreq 62%

basics

~20 s

A broadcast hash join ships one whole small table to every executor and probes it from inside each existing partition of the large table. No rows are redistributed by key, so a dominant key value never concentrates into a single task.

open as a page

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

level: middleimportance: should knowfreq 68%

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.

open as a page

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

level: seniorimportance: should knowfreq 56%

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.

open as a page

For a nightly Spark job with a recurring hot key, how do you choose between AQE, salting, and isolating the key?

level: principalimportance: should knowfreq 38%

basics

~20 s

Start with the free options: leave adaptive skew handling on and broadcast the small side if you can. Salt when a known hot key recurs and the operator is an aggregation or an unbroadcastable join. Isolate the key into its own path when one or two values dominate everything.

open as a page

In Spark AQE, how can a sort-merge join become a broadcast hash join mid-query?

level: seniorimportance: nice to knowfreq 28%

basics

~20 s

Adaptive Query Execution re-plans at shuffle boundaries. Once a join side's shuffle has been written, Spark knows its exact size; if that is under the adaptive broadcast threshold, it replaces the planned sort-merge join with a broadcast hash join.

open as a page