A step's 400 pieces average 9 seconds each, yet the step runs 52 minutes — what do the median and maximum piece durations add?
answer
- a step waits for its last piece
- the mean hides a tail
- median against maximum, not the average
- maximum close to the step's wall time
basics
~20 sThe 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 sA 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
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.
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.
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.
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.