A teammate says a running job 'feels slow' and gives no other detail. Which numbers do you read first, and why those?
answer
- read before you change
- rate, duration, distance, failures
- records per second is not bytes per second
- per step, not per run
- where and since when, before why
basics
~20 sFive 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 sI 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
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.
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.
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.
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.