skip to content

A nightly job's piece count was sized when its input was 200 GB and it now reads 3 TB unchanged - what rots, and how would you tell?

level: seniorimportance: should knowfreq 50%

answer

  1. the number froze, the data did not
  2. bytes per piece grew fifteenfold
  3. whole distribution shifted, no tail
  4. store a rule, not a literal

basics

~20 s

Bytes per piece grew fifteenfold while the count stayed put, so every piece is now oversized: uniformly longer runtimes, new disk spilling, memory pressure and expensive retries. The tell is that the whole duration distribution shifted, not that one piece is slow.

solid answer

~40 s

The count did not change, so the input growth landed entirely on the size of each piece. Every piece now holds roughly fifteen times the records: its working set no longer fits the memory one worker thread is given, so it **spills** - writes part of that working set to local disk and reads it back - and each lost piece costs fifteen times the recomputation. The diagnostic that separates this from other causes is the *shape* of the change: all pieces got slower by a similar factor. One piece far slower than the rest is either records concentrating on one key or one bad machine. The durable fix is to stop storing a literal number and derive it each run from current bytes and current lanes.

go deeper

for a junior

Recall that a fixed number of pieces over a growing input means each piece gets bigger, and that a bigger piece takes longer and needs more memory in the one thread that holds it.

for a middle

Explain the consequences precisely: working sets that no longer fit and start spilling to local disk, retries that cost proportionally more, and a wall clock that adding machines will not improve.

for a senior

Show the diagnosis by evidence shape - a shifted distribution rather than a tail - rule out records concentrating on a few keys and one slow machine explicitly, and replace the literal with a rule the next run recomputes.

for a principal

The angle is the fleet: a sizing policy expressed as a computation with a floor and a ceiling, per-job overrides that carry a reason, and drift surfaced by alerting on measured bytes per piece rather than discovered by a slow night.

## Why a right number stops being right The **piece count** is how many slices the input is divided into, each read and processed by one **worker thread** - one lane inside a worker process - from start to finish. It is a number chosen against two moving quantities: the bytes the step handles, and the lanes available to run it. Both move, and the number does not. Three separate drifts rot a count, and they have different signatures: 1. **The input grew.** Bytes per piece rise in lockstep. This is the case in the question: 200 GB over, say, 800 pieces was 256 MB a piece; 3 TB over the same 800 pieces is nearly 4 GB a piece. 2. **The cluster changed.** Lanes were added or taken away, so a count that gave clean waves now gives a ragged final one, or leaves capacity unused in every wave. 3. **The records changed shape.** The same byte count now decodes into more or heavier in-memory objects after a schema or format change, so memory per piece rises even though the stored size did not. How records are represented in memory is owned by Record Representation; the point here is only that bytes on disk are an imperfect proxy for the pressure on a lane. ## What oversized pieces look like from outside - **Every** piece duration rises by a similar factor - the whole distribution shifts right, rather than growing a tail. - Local disk write volume appears where there was none, because lanes are **spilling**: writing part of the working set to disk when it no longer fits the memory the lane is given. - Memory pressure and garbage-collection time rise across all lanes, not one. - A single lost piece costs far more than it used to, so an occasional machine failure now visibly moves the job's finish time. - The wall clock barely improves when machines are added, because the count, not the lane supply, is what bounds the run. ## Telling it apart from its two look-alikes | observation | oversized pieces | records concentrated on few keys (skew) | one slow machine (a straggler) | |---|---|---|---| | duration distribution | whole distribution shifted right | a long tail on an otherwise fast body | one outlier | | records per piece | uniformly high | wildly uneven | normal | | repeats on rerun | yes, every run | yes, and on the same keys | usually moves to a different machine | | who owns it | this decision | Skew and Stragglers | Skew and Stragglers | The distinction matters because the remedies do not overlap. A larger count fixes the first and does very little for the second, where the problem is the distribution of keys rather than the number of pieces. ## The durable form of the decision The defect is not the number that was chosen - it was right at 200 GB. The defect is that a *literal* was stored where a *rule* belonged. A durable sizing decision: - computes the count each run from the bytes the step will actually handle and the lanes currently available, rather than reading a value fixed a year ago; - carries a floor, so a small input is not shredded into pieces whose fixed per-piece cost exceeds their work; - carries a ceiling, so an unexpectedly large input does not create a count the coordinating process cannot track; - is reviewed on an event, not a calendar: the input crossing a size band, the cluster being resized, the format changing. Where the engine revises the count itself between step groups, this matters less - though that capability exists for some operators in some engines and largely not at all where operators are long-lived, and it is owned by Replanning at Runtime. ## Continuous jobs rot too, differently A **continuous job** - a run over an input with no end - has a **declared operator width** rather than a derived count, and it rots on arrival rate rather than on stored bytes. The traffic the pipeline was sized for doubles, each lane's share of the arrival rate doubles, and sustained throughput falls below the input rate; the job falls behind and the backlog grows. The difference is the cost of the correction: a finite job is corrected by resubmitting it tomorrow, while a width change here means a restart with per-key state redistributed over the new lanes. That asymmetry is the reason continuous pipelines deserve a headroom margin at the start that a nightly job does not. ## What an interviewer is checking That you diagnose from the shape of the evidence rather than from a hunch, that you can name the three look-alikes and say which one the evidence rules out, and that your fix outlives you: a rule the next run recomputes, not a bigger literal that will rot again at 30 TB.

  • How do you distinguish this from records concentrating on a few keys?
    By the shape of the duration distribution. Oversized pieces shift every piece rightwards by a similar factor and the record counts stay even. Concentration leaves most pieces fast and a handful enormously slower, with record counts to match. Raising the count helps the first and barely touches the second, which belongs to Skew and Stragglers.
  • The team wants one sizing policy across two hundred jobs. What makes that safe?
    Express it as a computation - bytes the step handles, divided by a target size, floored at the lane count and capped so the coordinating process stays healthy - rather than a shared literal. Let individual jobs override with a recorded reason, and alert when a job's measured bytes per piece leaves the intended band, so drift surfaces before a run degrades.
  • Does a count ever rot downward, so that the pieces become too small?
    Yes, in two ways. An upstream filter or aggregation can shrink what this step actually handles, leaving a count sized for the old volume producing tiny pieces. And a cluster that lost lanes makes a large count more waves rather than more parallelism. Both raise the fixed per-piece cost relative to the work done.

saying these in an interview costs you the question

  • Diagnosing uniformly slower pieces as records concentrating on one key
  • Fixing it by adding machines while the piece count stays fixed
  • Raising the count to a new literal that will rot the same way
  • Treating spilling to disk as a broken engine rather than a sizing signal
  • Assuming the engine notices the growth and re-sizes by itself
  • Reviewing sizing on a calendar rather than when the input crosses a band