skip to content

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%

answer

  1. duration alone is not evidence
  2. compare what the piece was handed
  3. bytes and records, against the median piece
  4. proportionate, or not proportionate

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.

solid answer

~50 s

Duration on its own only tells me a piece was slow. The second reading is what that piece was *handed*: bytes read and records read, compared with the median piece of the same step. There are two outcomes and they point in opposite directions. Either the piece read many times its peers' input, in which case the work was divided unevenly and the long duration is roughly proportionate to it; or it read about what everyone else read and still ran forty times longer, in which case volume does not explain the duration at all. A third reading is worth keeping in mind: equal bytes but far more records, which happens when records vary in size or compress differently, and then per-record cost dominates rather than volume. Which of those two causes is at work, and what to do about it, is the next question — but it cannot be asked honestly from duration alone.

go deeper

for a junior

Remember that a slow piece has two possible explanations and that the next number to fetch is how much input that piece read, compared with a typical piece in the same step.

for a middle

Explain the four combinations of bytes and records against the median, including why equal bytes with many more records points at per-record cost rather than at volume.

for a senior

Show that you read input size on the gathering side of the step that was slow, not on the source, and that you can still reach a verdict when the runtime reports only one of the two numbers.

for a principal

Consider what the team can actually see. If the runtime exposes no per-piece input size, the cheaper long-term fix is a standing measurement of the data rather than a deeper investigation of every slow run.

## Duration is the symptom, input size is the evidence A long piece is a fact with at least two explanations, and no amount of staring at the duration separates them. A **piece** is one share of a step's input that one **worker** — a process on one machine, with its own memory and local disk — handles independently. If one piece ran forty times longer than its peers, either it was handed roughly forty times the work, or it was handed an ordinary amount of work and took forty times as long to do it. The measurement that decides between those is the piece's **input size**, and it comes in two numbers that do not always agree: bytes read and records read. ## Where to read the input size This is the first place the question goes wrong, because two different objects are both called a piece: - **the input pieces the reader was handed** — how the source was cut up before anything ran. These are usually near-uniform by construction, because whoever cut them was aiming for uniform sizes. - **the pieces on the gathering side of a redistributing step** — a redistributing step being one that cannot be computed from what a single worker already holds, so every record is first sent to the worker that owns its key. The amount each of those pieces receives is decided by the data, not by the reader, and this is where imbalance almost always lives. So read the input size of the pieces *in the step that was slow*, not of the job's source. A perfectly even read of the input tells you nothing about what the gathering side received. ## The four readings | bytes for the long piece | records for the long piece | reading | |---|---|---| | many times the median | many times the median | the work was divided unevenly; the duration is roughly proportionate to what the piece was handed | | about the median | about the median | volume does not explain the duration — the cause is elsewhere | | about the median | many times the median | many small records; per-record cost dominates, not the bytes | | many times the median | about the median | few very large records; memory pressure and per-record size dominate | The middle two rows are the ones people skip, and they are the reason both numbers are worth reading. A step whose per-record cost is high — parsing, a nested structure, an expensive function per row — is sensitive to the record count in a way that byte totals disguise completely. ## Bytes and records disagree more often than people expect They diverge whenever record size is not constant, and it rarely is. Variable-length text fields, optional nested structures, and compression that works better on some pieces than on others all break the proportionality. A piece can read the median number of bytes and hold eight times the median number of records because its records happen to be small and highly compressible; the same piece will look ordinary on any byte-based check. Where a runtime reports only one of the two numbers, know which one you are looking at and what it cannot show you. ## Compare against the median, not the mean The outlier is inside the mean. With one piece at forty times the rest, the mean input size sits well above the typical piece and shrinks the apparent ratio, sometimes enough to make a real imbalance look like a mild one. The median is unmoved by a single extreme value and is the honest baseline. The same argument applies to the duration axis, which is why the two are read the same way: median as the reference, maximum as the question. ## When the runtime shows you less than this What a runtime exposes about each unit varies widely, and some publish duration far more readily than input size. Two fallbacks: 1. **Reason from proportion with whatever you have.** A forty-to-one duration against a one-to-one byte share is already conclusive, even without record counts. 2. **Measure the data instead of the run.** Counting records per key over the input, before the step, gives you the distribution the gathering side is going to receive. That is a separate pass, and often the better tool when the run report is thin. ## What the pair still does not settle Duration plus input size tells you whether the time is proportionate to the work. It does not name the cause: a piece that read an ordinary share and still ran long may be sitting on an unhealthy machine — a **straggler**, meaning slow for reasons of the machine rather than of the data — or may be running an expensive computation over ordinary input. And a piece that read forty times its peers' input has an uneven distribution behind it whose origin is a separate subject again. The pair narrows two possibilities to one; it is a checkpoint in the diagnosis, not the end of it.

  • The runtime reports bytes per piece but no record counts — is the check still worth doing?
    Yes. Bytes answer the main question — was this piece handed a proportionate share — and a forty-to-one duration on a one-to-one byte share is already conclusive. What you lose is the small-records case, where bytes match and record counts do not. You can usually recover that by counting records per key over the input separately.
  • Why compare the long piece against the median rather than the mean?
    Because the outlier is in the mean. One piece at forty times the rest pulls the mean well above the typical piece and shrinks the apparent ratio, which can make a severe imbalance read as a mild one. The median ignores a single extreme value and gives the honest baseline.

saying these in an interview costs you the question

  • Declares data skew from a long piece duration without reading its input size.
  • Assumes bytes read and record counts always move together.
  • Compares the long piece against the mean piece rather than the median.
  • Reads the input size of the whole step instead of per piece.
  • Checks the source's input pieces when the imbalance is on the gathering side.