Four hundred pieces finish in seconds and one runs an hour: what evidence tells you whether the data or the machine caused it?
answer
- two columns: bytes read, time taken
- compare against the middle, not the average
- does it follow the host or the unit
- same key, same destination, every run
- no finished units means compare live instances
basics
~20 sCompare 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.
solid answer
~50 sStart with the pair of numbers the runtime records per unit and puts in its run report: **how much input each one read** and **how long it took**. If the long-running unit read roughly what its peers read, the data is not the cause and the machine is - that is a straggler. If it read forty times as much, the cause is data skew, an uneven number of records per piece. Two confirmations settle the remaining doubt: whether other units placed on the same machine were also slow, and whether a re-run puts the slowness on the *same* unit (the key-to-destination rule is stable, so a data cause reproduces) or on a *different* unit (the slowness travelled with the host). Only then choose a remedy: rewriting the job addresses one cause, a second copy of the unit addresses the other, and each is useless against the wrong one.
go deeper
Remember that a long duration on its own proves nothing; you need the input size beside it before you can say which cause you are looking at.
Explain why the outlier is compared with the middle of the pack rather than the average, and why the same key always lands on the same destination.
Walk the evidence in order - sizes, then other units on the same host, then a re-run - and state explicitly which remedy each cause admits and why the other is wasted.
Ask what it costs the organisation to keep diagnosing this by hand, and whether the money goes to better per-unit reporting, to duplicated work by default, or to fixing the key upstream.
## The one measurement that does most of the work Every runtime records, for each unit of a step, roughly how much input it read and how long it ran, and exposes them in some run report - whatever view it offers of what each unit read and how long it took, in flight or once the step is done. Those two numbers together answer the question; either one alone does not. | What the outlier read | What it means | Name | |---|---|---| | Far more than its peers | It is doing genuinely more work at a normal rate | Data skew - an uneven number of records per piece | | About the same as its peers | It is doing normal work at an abnormal rate | A straggler - the machine under it is the problem | The comparison must be against the *spread* of the other units, not against an average of all of them: an average that includes the outlier is dragged toward it and hides exactly the gap you are looking for. Compare the outlier with the middle of the pack. ## Three confirmations when the first reading is ambiguous Input sizes are not always available or trustworthy - compressed sources, wildly variable record widths, and a join that multiplies rows can all make "bytes read" a poor proxy for work done. Three further tests are cheap: 1. **Look at the other units that ran on the same machine.** If several unrelated units placed there were also slower than their peers, the host is the common factor and no property of your key explains it. 2. **Run it again and see what travels.** The rule that maps a record's key to its destination is stable, so a given key lands on the same destination number run after run. If the *same* unit number is slow again, on different machines, the cause is in the data. If a *different* unit is slow and it happens to be on the same host as last time, the cause is the host. 3. **Count records per key on the input.** If one value covers a large share of the rows, you have an explanation for a heavy unit before you look at any machine at all. ## What varies by runtime, and why it matters here This whole procedure assumes there are finished units with recorded sizes to compare. That assumption is a property of the model, not of the world: - In a finite job made of steps that complete, the report of finished units is exactly what you need, and the comparison is over a population of peers that have all ended. - In a **continuous job over an endless input** - one whose input never ends, so no step ever finishes - there are no completed units to rank. What you compare instead is the live parallel instances at the same moment: their throughput, and their **backlog**, meaning how far behind the newest available record each one currently is. The reasoning is identical in shape; only the evidence differs. An instance whose backlog grows while its input rate matches its peers' points at its machine; one that receives several times the records of its peers points at the key. - In runtimes where a redistributing step - a step that cannot be computed from what one worker already holds, so every record is first sent to the worker that owns its key - writes its output to local disk before the next step reads it, disk health on one host is a much larger share of the story than in a design that pushes records straight across the network. ## The remedies do not transfer This is why the diagnosis is worth the ten minutes: - **Machine-caused.** A second copy of the unit on a healthy machine helps, where the runtime offers one. Not placing work on that machine helps, though whose surface that is depends on how the job gets its machines. Nothing about the job's logic needs to change. - **Data-caused.** No amount of duplicated work helps: the copy holds the same records. The fix is in the job or in the data upstream, and it costs author time and a correctness argument. And when both are true at once - a somewhat heavier unit that also landed on a sick machine - deal with the data cause first. It is deterministic and it reproduces, so you can verify the fix; once the shares are even, the machine's contribution shrinks from a tenfold outlier to a small multiple that duplicated work can absorb. ## The answer that fails "One task is slow, so there is skew, so I will spread the key" is the answer interviewers are listening for and marking down. It skips the measurement, commits author time to a rewrite, changes the shape of the job, and leaves a failing disk exactly where it was. The mirror-image failure - "turn on duplicate work and move on" - burns a second machine for the same hour whenever the unit is genuinely heavy. The evidence costs one look at two columns.
- The unit is both heavier than its peers and on a machine that is misbehaving. Where do you start?With the data cause. It is deterministic, it reproduces on every run, and it is the one duplicated work cannot touch. Once the shares are even you can measure again: what is left is the machine's contribution, now a small multiple rather than a tenfold outlier, and that is the part a second copy of the unit can absorb.
- Why is input size read sometimes a poor proxy for how much work a unit has?Because bytes are not records and records are not equal work. Compressed input expands by wildly different factors, one record can carry a large nested structure, and a join can multiply a key's rows so that a modest input produces an enormous intermediate. Where that is a risk, count records per key on the input rather than trusting bytes.
- How do you run the same check on a job whose input never ends?There are no finished units to rank, so compare live parallel instances at the same moment: their throughput and their backlog. An instance whose backlog grows while it receives a normal share of records points at its machine; one receiving several times its peers' records points at the key.
saying these in an interview costs you the question
- Declares skew from duration alone without checking input size
- Compares the outlier against the average of all units
- Thinks a re-run reshuffles which key lands where
- Proposes a key rewrite before identifying the cause
- Assumes there is always a report of finished units to read