In a Spark stage, which task metrics separate data skew from a slow executor?
answer
- compare tasks, not the whole stage
- look at the max-versus-median row
- is the straggler reading more, or just slower?
- one host, or one partition?
- group the key and count the rows
basics
~20 sCompare 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 sOpen 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 linesSummary 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 GBgo deeper
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.
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.
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.
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