skip to content

Throughput, Lag and Failures

The handful of numbers that say a running job is healthy — work done per second, time per step, distance behind the input, failed units — and which moves first when it is not.

on this pageshow

questions

5

A teammate says a running job 'feels slow' and gives no other detail. Which numbers do you read first, and why those?

level: juniorimportance: must knowfreq 72%

answer

  1. read before you change
  2. rate, duration, distance, failures
  3. records per second is not bytes per second
  4. per step, not per run
  5. where and since when, before why

basics

~20 s

Five numbers cover almost every 'slow' report: processed throughput in records and in bytes per second, duration per step of the job graph, distance behind the input, and failed or retried units of work. Read them before changing anything.

solid answer

~50 s

I read five numbers before touching a setting. **Processed throughput**, as records per second and bytes per second separately, because they move independently. **Duration per step of the job graph** — a step being one transformation applied to every piece of the input — so I know which transformation grew rather than 'the job'. **Distance behind the input**: lag in the sense of unread input, how far the newest record the job has finished with sits behind the newest record its input already holds. That figure only exists where the input keeps growing. **Failed and retried units of work**, a unit of work being the smallest thing the engine hands a worker process and the smallest thing it retries by itself. And **memory headroom with bytes spilled to a worker's local disk**. Together these say where and since when; nothing is worth changing until they do.

go deeper

for a junior

Be able to name the numbers and say what each measures: work done per second in records and in bytes, time per step, distance behind the input, failed or retried units, bytes spilled.

for a middle

Explain why records per second and bytes per second must be read as two figures, and why a per-step duration attributes a slowdown that a whole-run duration cannot.

for a senior

Show the order and defend it: rate, then location, then impact, then the failure counters that predict the next hour. Say which figures you had to arrange in advance to have at all.

for a principal

Argue what the team deliberately does not measure and why, and what it costs when the only people who can read these numbers are the three who built the job.

## What a 'slow' report actually contains 'Slow' is a feeling, not a measurement. Before it becomes a question anyone can answer, three things have to be attached to it: **which part** of the run got slower, **by how much** against a normal run, and **since when**. A *job run* here means one submission of one program to the cluster, from start to its last written output. The numbers below exist to supply exactly those three things, and none of them requires you to know what the program computes. The habit that separates an engineer who can operate a job from one who cannot is reading before changing. Adding machines, raising a memory allowance or increasing the number of pieces the input is cut into are all *changes*, and a change made before a measurement destroys the only comparison you had. ## The five numbers 1. **Processed throughput, as two figures.** Records per second and bytes per second, read separately. They move independently: the record rate can sit flat while bytes per second doubles because the records themselves grew, and the run then dies of volume while the record rate reports nothing unusual. 2. **Duration per step of the job graph.** A *step* is one transformation in the job's graph, applied to every piece of the input. A whole-run duration tells you the run is slower; a per-step duration tells you which single transformation is slower, and that is usually most of the diagnosis. 3. **Distance behind the input** — lag in the sense of unread input: how far the newest record the job has finished with sits behind the newest record its input already holds, in records or in seconds. This figure exists only for a job reading a source that keeps growing. A finite input has an end rather than a growing frontier, so there is nothing to be behind. 4. **Failed and retried units of work.** A *unit of work* is the smallest thing the engine hands to a **worker process** — one process on one machine that runs pieces of the job and owns the memory they use — and the smallest thing the engine will retry by itself. Those retries usually do not surface outside the run, so the count is the only evidence they happened at all. 5. **Memory headroom and bytes spilled.** *Spill* means bytes a worker writes to its own local disk because a working set will not fit in memory. On this reading exercise it is simply a number that has started rising; the memory budget behind it, and what to do about it, belong to the memory subject rather than to this one. ## What each one rules out | Number that moved | Rules in | Rules out | |---|---|---| | bytes/s up, records/s flat | records grew larger; more data per piece | a change in your own program's logic | | one step's duration up, the rest flat | that transformation, its volume, its memory | cluster-wide slowness | | every step's duration up | contention, fewer machines, a shared pool | your program | | distance behind rising | a sustained shortfall between work done and work arriving | a single one-off blip | | retried units up | a sick machine, or a piece of input that keeps failing | a purely arithmetic slowdown | | spilled bytes up | a working set that stopped fitting in memory | a network problem | ## Where the execution model changes the reading This class of engine does not share one runtime, and two of the five figures mean different things depending on the model in front of you. - On a **continuous job built from repeated small finite runs** — the engine slices an endless input into small finite jobs run back to back — durations arrive one slice at a time, and the latency floor is one slice's period. The signal is a slice taking longer than its own period, not the raw duration. - On a **record-at-a-time runtime with key-bound state**, where each record moves through the graph as it arrives and is folded into state held under its key, there are no run boundaries inside the job to time. The equivalent of a per-step duration is the fraction of its time each step spends busy, read continuously. - On the **two-phase disk-to-disk model**, where every phase writes its whole output to shared storage before the next phase reads it, durations come naturally per phase, and one phase may contain several graph steps — so 'the slow step' is a coarser answer there than on the other two. Whatever screen or endpoint a given engine exposes its live and finished-run numbers on, refer to it by what it shows rather than by its name. The names differ between engines; these quantities do not. ## The order is the answer An interviewer asking this is not collecting a list, they are checking whether you have an order and a reason for it. A defensible order runs: rate first (is work still being done at all), then per-step duration (where), then distance behind the input (how bad, and for whom), then failures and spilled bytes (why it may still be getting worse). Answering 'I would add workers' before any of that is precisely the reflex the question exists to catch.

  • Which of these numbers exists only because someone arranged for it in advance?
    Throughput and per-step duration are usually published by the runtime itself. Distance behind the input often needs the source's frontier compared against the job's position, which someone has to wire up. A per-record or per-key count you care about is never free: it has to be incremented by the workers and summed by the coordinating process so the evidence outlives the run.
  • Why is the whole run's wall-clock duration a poor first number?
    It confirms the complaint and attributes nothing. Every cause — a larger input, one slow transformation, fewer machines, a piece of work retried repeatedly — produces the same longer duration. It is useful only as a comparison against previous runs of the same shape, and only to decide whether to investigate at all.
  • The job reads a finite input, so there is no distance-behind figure. What replaces it?
    Fraction of the input consumed against elapsed time, which gives a projected finish, plus the count of pieces completed against total pieces. Those answer the question distance-behind answers for a growing source: how far through are we, and will we arrive when we promised.

saying these in an interview costs you the question

  • Proposes adding machines before reading a single number.
  • Treats records per second as the only throughput figure.
  • Quotes the whole run's duration and never a per-step figure.
  • Assumes every job has a distance-behind-the-input number, including finite ones.
  • Calls a run healthy because the process is still alive.
  • Reads only the latest value and never a comparison with a normal run.
open as a page

Which signals on a running job move first when it degrades, and which only confirm the damage afterwards?

level: middleimportance: must knowfreq 58%

basics

~20 s

Leading signals move while output still looks fine: bytes spilled to local disk, shrinking memory headroom, retried units of work, and the fraction of a step's time spent unable to pass its output on. Distance behind the input and run duration are lagging — they accumulate the damage.

open as a page

Records per second is unchanged but each step of the job now takes twice as long — which further measurements separate the causes?

level: middleimportance: should knowfreq 55%

basics

~20 s

Read bytes per second beside records per second, plus bytes per record. A flat record rate with doubled step duration usually means fatter records, more input per piece, or slower machines — three different fixes, separated by the volume figures and by the spread across pieces.

open as a page

Every health check is green and the process is up, yet consumers say the numbers stopped updating hours ago — what should the alarm have measured?

level: seniorimportance: should knowfreq 52%

basics

~20 s

Liveness is not progress. The useful alarm is a threshold on what the job emits — the freshness and volume of its own output, or its distance behind the input — evaluated by something outside the job, so that silence itself trips the alarm.

open as a page

A job run finished successfully but retried four hundred times as many units of work as usual — why investigate?

level: seniorimportance: should knowfreq 44%

basics

~20 s

A green run hides its retries: the engine re-ran failed units until they passed, so success proves only that the allowance was large enough. The count names a machine, an input piece or a memory limit that is failing today and will exhaust that allowance soon.

open as a page