A piece reading a normal share of input runs ten times longer than its peers - what is that called and why?
answer
- same input, very different duration
- the machine, not the key
- failing disk, busy neighbour, cold start
- slowest participant sets the pace
basics
~20 sThat 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.
solid answer
~40 sA job is cut into `pieces` - one share of a step's input that one worker processes on its own - and those pieces are meant to be interchangeable, so a run costs roughly what its slowest piece costs. When one piece runs long there are two possible causes. Either it holds far more records than its peers (data skew - an uneven number of records per piece), or it holds an ordinary share and the machine underneath it is unhealthy: a disk failing and retrying reads, another process on the same host taking CPU or network, a machine that started cold and has nothing cached yet. The second case is the **straggler**, and the giveaway is exactly the pattern in the question: comparable input read, wildly incomparable duration.
go deeper
Be able to say that a unit reading a normal share of input but running far longer points at its machine, and name two or three machine causes such as a failing disk, a busy neighbour or a cold start.
Explain why parallel work costs what its slowest participant costs, and why fresh capacity does nothing for a unit that has already been placed on an unhealthy machine.
Show that you check input size beside duration before proposing anything, and that you know a rewrite of the grouping key cannot repair hardware, however slow the unit looks.
Treat it as a fleet question: how often machine-slow units occur across all jobs, whether paying for duplicate work is worth it, and when the honest answer is to stop placing work on bad machines.
## The two reasons one piece can run long A distributed job is cut into **pieces** - one share of a step's input that one worker handles independently of the others - and a **worker** is one process on one machine, with its own memory and its own local disk, that runs those pieces. The pieces are meant to be interchangeable, so the arithmetic of the whole thing is simple: parallel work finishes when its slowest participant finishes. Four hundred pieces that each take twenty seconds cost twenty seconds; three hundred and ninety-nine of those plus one that takes an hour costs an hour, with the cluster mostly idle for fifty-nine minutes of it. When one piece runs long, exactly one of two things is true, and the entire practical value of this subject is telling them apart: - **The piece holds far more work.** This is *data skew - an uneven number of records per piece*. The grouping key that piece owns covers a large share of all the rows, or the source was cut unevenly. The piece is doing genuinely more, and it is doing it at a normal rate. - **The piece holds an ordinary share and its machine is not healthy.** This is a **straggler**. The piece is doing the same amount of work as its peers and taking much longer to do it. The cheapest evidence is a pair of numbers the runtime already records for every unit and shows in the run report - whatever view it offers of what each unit read and how long it took. **Comparable input with incomparable duration is the straggler's signature.** A far larger input read is the other cause entirely. ## What actually makes a machine slow None of these have anything to do with your data or your key: - **A degrading disk.** Reads that fail and retry, or a device whose throughput has quietly collapsed, slow every byte the piece touches - and a unit that writes intermediate results to local disk touches a lot of them. - **A noisy neighbour.** Another process on the same host - another job's worker, a log shipper, a backup - taking CPU time, memory bandwidth, disk queue depth or network capacity that your worker assumed it had. - **A cold start.** A machine that joined recently has nothing resident yet: no bytes fetched from remote storage, no pages warm, no runtime code warmed up. The first pieces it runs pay costs its peers already paid. - **A degraded network path.** One host behind a saturated or flapping link fetches its input, or sends its output to other workers, at a fraction of everyone else's rate. - **An oversubscribed host.** More worker processes placed on the machine than its cores or its memory bandwidth can actually serve, so each one runs at a fraction of speed. ## Why one slow piece costs the whole run How the delay surfaces depends on the runtime model, and this is where a candidate who has only used one engine usually overclaims. Where a step has a boundary - the next step cannot start until every piece of this one has produced its output, because records have to be moved between workers first - the step genuinely ends when its last piece ends, and the slow piece sets the wall-clock cost directly. Where the runtime is pipelined and record-at-a-time over an input that never ends, no step ever finishes at all; the slow instance instead accumulates a **backlog** of records it has not reached, and in most such designs the pressure travels back to whatever is feeding it. In both models the same thing is true in the end: the slowest participant sets the pace, and capacity elsewhere cannot spend itself on the piece that is behind. ## The contrast at a glance | | Heavy piece (the data) | Straggler (the machine) | |---|---|---| | Input read | Far more than its peers | About the same as its peers | | Reproducible | Yes - follows the key, so the same piece is slow again | No - follows the host, so a different piece is slow next run | | Other units on that host | Normal | Often slow too | | What helps | Rewriting the job so the work is spread | A second copy on a healthy machine, where the runtime offers one | | More machines alone | No help | No help | ## Why naming the cause comes before the fix The expensive mistake is reaching for a key rewrite the moment a long-running piece appears. Rewrites of that kind cost author hours, have to be re-reasoned for correctness, and do absolutely nothing to a disk that is retrying its reads. The mirror-image mistake also exists: throwing a second copy of the unit at a piece that genuinely holds forty times the records just burns a second machine for the same hour. One measurement - input read beside duration, per unit, against the peers in the same step - decides which family of remedy is even applicable, and it is available before you change a line of the job.
- Does adding more workers to the cluster shorten a step held up by one unhealthy machine?No. The slow piece is already assigned to a worker and arriving capacity cannot divide it or take it over. New machines can only pick up work that has not started. The only ways the extra capacity gets used against that piece are a second copy of it, where the runtime offers one, or the piece failing and being placed again.
- Why does a machine that only just joined run the same work more slowly than its peers?Because nothing it needs is resident yet. Input bytes have not been fetched from remote storage, nothing is warm in page cache, and the runtime's own code has not been through whatever warm-up it does. Its peers paid those costs on their first pieces; this machine pays them now, which shows up as a slower unit at a normal input size.
saying these in an interview costs you the question
- Calls every long-running piece skew and immediately rewrites the key
- Assumes adding machines shortens a piece that is already placed
- Reads the average unit duration instead of the spread
- Thinks straggler means the unit was handed more records
- Believes a slow unit always means the job was written badly