skip to content

Skew and Stragglers

Uneven work: one piece holding far more records than the rest (skew), or one worker running slowly for reasons of its own (a straggler), and every remedy for either, whoever applies it.

on this pageshow

explore

questions

20

A step's 400 pieces average 9 seconds each, yet the step runs 52 minutes — what do the median and maximum piece durations add?

level: juniorimportance: must knowfreq 76%

answer

  1. a step waits for its last piece
  2. the mean hides a tail
  3. median against maximum, not the average
  4. maximum close to the step's wall time

basics

~20 s

The median shows the typical piece finished in about a second while the maximum shows one ran nearly the full 52 minutes: a long tail of one or two outsized pieces, not a step that is uniformly slow.

solid answer

~50 s

A step is not finished until its last **piece** is finished — a piece being one share of the step's input that one worker handles independently — so the step's wall time is set by the *maximum* piece duration, not by the mean. The mean of 9 seconds tells me how much machine time the step consumed in total (`400 x 9`, about 3,600 seconds); it is a cost figure, not a latency one. What I want is the median, around a second here, and the maximum, which will be close to the 52-minute wall time. That pair is the signature of a long tail: almost everything finished immediately and one piece carried the step. A uniformly slow step looks different — every piece within a small factor of the median. The tail case will not shorten just because I add machines.

go deeper

for a junior

Recall that a step ends when its slowest piece ends, so the number to hunt for is the longest piece duration, with the median beside it as the point of comparison.

for a middle

Explain what the mean, the median and the maximum each measure, and use the arithmetic to show why a small mean next to a long wall time can only mean a concentrated tail.

for a senior

Demonstrate that you read the distribution before forming a theory, pair duration with input size per piece, and refuse to buy machines for a step whose cost sits inside one piece.

for a principal

Weigh the cost of the diagnosis itself. A team whose dashboards report step averages will keep enlarging clusters, so the per-piece distribution has to be the figure that reaches whoever approves the spend.

## The number a step is judged by A step of a distributed job is cut into **pieces** — one share of the step's input that one **worker** (one process on one machine, with its own memory and its own local disk) handles independently of the others. Pieces run concurrently, as many at a time as there are free worker threads, and **the step is not finished until its last piece is finished**. That single sentence decides which statistic is worth reading. The step's wall time is set by the longest piece; the mean piece duration is a statement about how much machine time the step consumed, which is a cost question rather than a latency one. ## Running the arithmetic in the example Four hundred pieces at a mean of 9 seconds means the step bought roughly `400 x 9 = 3,600` seconds of machine time. The step ran for 52 minutes, or 3,120 seconds of wall time. Had that machine time been spread evenly over, say, forty concurrent worker threads, it would have come back in about 90 seconds. It did not, so nearly all of those 3,600 seconds sat inside one piece: - **median** piece duration: about 1 second — the typical piece - **maximum** piece duration: about 51 minutes — very close to the step's own wall time - the other 399 pieces together: a few hundred seconds The mean of 9 seconds is an artefact of dividing one enormous piece across four hundred of them. It describes no piece that actually ran. ## The two shapes, and what each one means | per-piece durations | reading | do more machines shorten it? | |---|---|---| | all within a small factor of the median, all long | the step is uniformly slow — every piece has real work to do | yes, roughly in proportion to the extra concurrency | | median small, maximum near the step's wall time | a long tail — the cost is concentrated in one or a few pieces | no; the extra machines sit idle waiting for that piece | | a broad middle band with a handful several times the median | a partial tail — worth measuring, often not worth acting on | partly, and less than the piece count suggests | That is the whole diagnostic value of the spread, and it is a lot. It does not yet say **why** the long piece was long: an uneven share of records (**data skew** — an uneven number of records per piece) and an unhealthy machine underneath one worker produce a similar duration profile, and telling them apart needs a second reading. But it settles the question teams usually get wrong first — whether this is a job that needs more capacity, or a job whose capacity is already idle. ## Where the numbers come from Every runtime of this class offers some **run report**: whatever view it gives of what each unit read and how long it took, in flight or once the step is done. What that view contains varies a great deal between runtimes — some publish a per-piece table with duration, bytes and record counts; some publish only summary quantiles for the step; some publish nothing until you collect it yourself. Do not build a diagnosis that assumes a particular view. The minimum you need is a duration per piece, or failing that a median and a maximum. Two practical notes: 1. **You can read the spread before the step ends.** Completed pieces already have durations, and the ones still running have an elapsed time you can hold against them. Five pieces still running twenty minutes after the other 395 completed is the same signature, spotted forty minutes earlier. 2. **A step's wall time can exceed its longest piece.** When there are more pieces than free worker threads, the pieces run in waves and the step costs roughly the sum of those waves; queueing behind other work on a shared cluster adds to it. That is why the arithmetic above says *had it been spread evenly* rather than predicting a firm number. ## Why the mean survives as a habit Because it is the number a summary line reports. The step took an hour; it had four hundred pieces; that is nine seconds each. Every one of those statements is true, and together they imply a story — everything is a little slow — that the distribution flatly contradicts. The discipline is small and mechanical: never accept a mean without the median and the maximum beside it, and read the maximum against the step's own wall time rather than against a ratio somebody quoted. ## What to read next Duration is one axis. The second is **what each piece was handed** — bytes read and records read, compared across the same set of pieces. A piece that ran forty times longer and read forty times the input is doing proportionate work; a piece that ran forty times longer on an ordinary share of input is a different problem. Reading the two axes together is what turns "one piece is slow" into something you can act on.

  • Why can a step's wall time be longer than its longest single piece?
    Because pieces run in waves when there are more of them than there are free worker threads. With four hundred pieces and forty threads, ten waves of one-second pieces already cost ten seconds before the long piece is accounted for. Waiting behind other work on a shared cluster adds to it in the same way.
  • The step has not finished yet — can you read the spread already?
    Usually. The pieces that completed already have durations, and the ones still running have an elapsed time you can compare against them; a handful still running long after the rest completed is the same signature. How much a runtime exposes mid-flight varies, so some let you watch it live and others only show the picture once the step is over.
  • Is a maximum-to-median ratio of ten always a problem?
    No, and there is no published threshold worth quoting. What matters is whether the longest piece is close to the step's own wall time. Ten-to-one between a one-second median and a ten-second maximum costs nothing; the same ratio between five minutes and fifty minutes is the entire run.

Four hundred supermarket checkouts. The mean basket takes ninety seconds, but the store cannot lock its doors until the last customer leaves, and one lane has a trolley that will take an hour. The mean basket is a fact about how much the store sold. The longest queue is the fact about when everyone gets to go home.

saying these in an interview costs you the question

  • Reads the average piece duration and calls the whole step uniformly slow.
  • Expects more workers to shorten a step whose cost sits in one piece.
  • Treats the sum of all piece durations as the step's wall time.
  • Judges the spread from duration alone, never checking input size per piece.
  • Assumes a high maximum always means an unhealthy machine.
open as a page

In a step that gathers every record with one key onto a single worker, why does doubling the cluster not shorten the commonest key's work?

level: juniorimportance: must knowfreq 82%

basics

~20 s

A key-to-destination rule sends every record carrying the same key to one destination, and that mapping is fixed. Extra machines add destinations, not a way to split one key, so the commonest key still runs on a single worker.

open as a page

How does appending an artificial suffix to a heavy grouping key spread its records, and what second pass does that force?

level: juniorimportance: must knowfreq 60%

basics

~20 s

Appending a small integer to the grouping key turns one key value into several, so its records reach several destinations instead of one. A second pass then re-groups those partial results under the original key to produce the real answer.

open as a page

A piece reading a normal share of input runs ten times longer than its peers - what is that called and why?

level: juniorimportance: must knowfreq 58%

basics

~20 s

That is a straggler: a unit whose slowness comes from the machine under it - a failing disk, a busy neighbour process, an empty local cache - rather than from holding more records than its peers.

open as a page

A step finished and reported 200 MB where the plan assumed 200 GB - what can a runtime that re-plans mid-run change now?

level: middleimportance: must knowfreq 55%

basics

~20 s

Runtime replanning - re-deciding unrun steps from the sizes a finished step actually produced - can change only what has not started: the shape of the next step's pieces, or a join switched to copying the small side. Finished work stands.

open as a page

Why does grouping on a column where missing values were written as an empty string concentrate work on one worker?

level: middleimportance: must knowfreq 66%

basics

~20 s

Placeholders compare equal to one another, so every row that was missing a value forms a single enormous group. That group is one key, one key resolves to one destination, and one worker inherits all of it.

open as a page

Which aggregates return exactly the same answer after a heavy key is spread over suffixes and recombined, and which do not?

level: middleimportance: must knowfreq 52%

basics

~20 s

Aggregates whose partials combine with the same operation - sum, count, minimum, maximum - come back identical. A mean survives only if each partial carries a sum and a count. Exact distinct counts and percentiles do not combine that way at all.

open as a page

Four hundred pieces finish in seconds and one runs an hour: what evidence tells you whether the data or the machine caused it?

level: seniorimportance: must knowfreq 55%

basics

~20 s

Compare input read against duration per unit. An ordinary input with a huge duration means the machine; a far larger input means an uneven share of records. Confirm it by whether the slowness follows the host or the same unit across runs.

open as a page

After a redistributing step finishes and reports its actual output sizes, which changes can a runtime still make to the unrun grouping or join?

level: middleimportance: should knowfreq 45%

basics

~20 s

Commonly three result-preserving moves, where an engine offers them: hand several undersized destinations to one worker, divide one that came out far too large, and switch a join to copying the small side once that side measures small.

open as a page

Before a step that sends every record to the worker owning its key, how do you check whether the grouping key's values are evenly spread?

level: middleimportance: should knowfreq 52%

basics

~20 s

Count records per key over the input, using the exact key expression and filters the job will apply, then compare the heaviest key's count with the total record count and with the typical key's count.

open as a page

One piece of a step runs 40 times longer than its peers — which second per-piece measurement do you read next, and why?

level: middleimportance: should knowfreq 58%

basics

~20 s

Input size for that piece — bytes read and records read — against the median piece. Either it was handed a far larger share of the work, or it read an ordinary share and something other than volume explains the time.

open as a page

When a runtime starts a second copy of a unit running far behind its peers, what does that buy and what does it cost?

level: middleimportance: should knowfreq 46%

basics

~20 s

A duplicate attempt starts a second copy of a lagging unit on another machine and keeps whichever copy finishes first, discarding the other. It buys back machine-caused slowness, and it costs a second machine's time plus a second read of the same input.

open as a page

Runtime re-planning divided an oversized destination piece in a join, yet left an equally oversized one in a grouping alone - what differs?

level: seniorimportance: should knowfreq 40%

basics

~20 s

A join can be cut inside one key, because the other side's matching rows are copied to each sub-piece and the union is the same result. A grouping cannot: every record for one key must meet in one place.

open as a page

In a continuous job over an endless input, what stands in for per-piece duration when you are looking for uneven work?

level: seniorimportance: should knowfreq 44%

basics

~20 s

Per-piece backlog — how far behind the newest record each piece is — read as a spread across pieces. A few pieces falling further behind while the rest sit near zero is uneven work; all of them drifting together is not.

open as a page

In an equality join, one key value has two million rows on the left and fifty on the right — what must that key's destination produce?

level: seniorimportance: should knowfreq 58%

basics

~20 s

Two million times fifty is one hundred million output rows, all produced by the single worker that owns that key. A heavy key in a join does not only concentrate input on one destination — it multiplies that input into output there.

open as a page

You spread a heavy grouping key over 200 artificial suffix values and attach the suffix to every key - what does the width buy, and what does attaching it everywhere cost?

level: seniorimportance: should knowfreq 45%

basics

~20 s

Width divides the heavy key's records by that factor and buys nothing once its share drops below an ordinary piece's. Attaching the suffix to every key multiplies the intermediate result by 200, and the second pass must move and read all of it.

open as a page

A unit's function writes a row into an external system - why does that make a duplicate attempt on it unsafe?

level: seniorimportance: should knowfreq 40%

basics

~20 s

Because the runtime can discard the losing copy's records but not what that copy already did outside the job. Two copies means two rows in the external system, and dropping one copy's output does not remove them.

open as a page

Why does grouping a billion rows by a two-value status column leave most of a hundred-worker cluster idle?

level: middleimportance: nice to knowfreq 44%

basics

~20 s

The number of distinct key values is a ceiling on how many destinations can receive anything. Two distinct values means two destinations do all the work, however many workers the cluster has and however many pieces the step runs with.

open as a page

You spread a heavy join key on one input with an artificial suffix - what must the other input do, and what does that cost?

level: seniorimportance: nice to knowfreq 34%

basics

~20 s

Every row on the other input under a spread key must be copied once per suffix value, so each pairing can still meet. That multiplies those rows by the width, and the combining pass is a union rather than a re-aggregation.

open as a page

Your batch jobs lean on runtime re-planning around uneven pieces, and the same logic must now run continuously over an endless input - what coverage is lost, and how do you plan around it?

level: principalimportance: nice to knowfreq 28%

basics

~20 s

Nearly all of it. Re-planning needs a finished step's measured output, and an endless input has none. Uneven-work decisions move to design time and to the author, and changing one means a restart rather than a mid-run correction.

open as a page