skip to content

Two daily runs write outputs of the same size and one costs ten times the other, so which per-run numbers separate the causes?

level: seniorimportance: should knowfreq 48%

answer

  1. same output, very different bill
  2. four numbers per run
  3. bytes read against bytes needed
  4. healthy signals, tenfold invoice
  5. alarm on the trend, not the level

basics

~20 s

Four numbers per run: input bytes read, capacity held times duration, bytes moved between workers, and output bytes written. With equal outputs, a tenfold gap almost always sits in input read or in capacity held while idle.

solid answer

~50 s

Equal output is what makes the two comparable, not what predicts the bill, so compare inputs and effort instead. Record four numbers for each run: **input bytes read**, **capacity held times duration**, **bytes moved between workers**, and **output bytes written**. Equal outputs with a tenfold difference in money almost always means one run read far more input than its result depends on, or held a large cluster idle through a long tail of a few unfinished pieces. The useful derived figure is the ratio of bytes read to the bytes the result actually depends on. None of this shows up on operating signals: throughput, distance behind the input and failure counts are all green on a run reading forty times too much, because reading too much is not a malfunction. It is only visible against the input the job needed.

go deeper

for a junior

Recall that a run's money cost tracks what it read and how long it held machines, not the size of what it produced. A job can write one small file and still have read everything ever stored.

for a middle

Explain which four numbers distinguish the causes and why each is the dominant term for a different one: input read, capacity held times duration, bytes moved between workers, output written.

for a senior

Show that you captured those numbers per run in advance and compare a ratio over time. Name why health signals are silent here, and reach for the bound that stopped narrowing before blaming the engine.

for a principal

Make the case that reading too much, not runtime tuning, is what dominates a real bill, and decide what the organisation will routinely review. A step-change alarm on one ratio per job beats a tuning effort nobody repeats.

## Why the expensive run looks healthy The signals a running job publishes answer one question: is it working? Throughput in records and bytes, how far behind the input it is, how many units of work failed and were retried, how much a worker wrote to local disk when its working set would not fit in memory. A run that reads forty times more input than its result depends on scores perfectly on every one of them. It is not failing, not lagging, not retrying. It is doing an enormous amount of work correctly, and the only surface on which that is visible is the invoice — which arrives weeks later and, in a shared pool, names the pool. This is the gap that makes a cost review a distinct exercise rather than a by-product of monitoring. ## The four numbers For every run, keep: - **Input bytes read** — the volume actually fetched from shared storage, meaning the storage every machine in the cluster can read, as opposed to a worker's own local disk which vanishes with the worker. - **Capacity held times duration** — worker processes multiplied by the seconds they were held, busy or not. - **Bytes moved between workers** — how much had to be redistributed across the network for steps that need records held elsewhere. - **Output bytes written** — and the number of files, because a destination charges per operation as well as per byte. All four are needed because each is the dominant term for a different failure. Reading too much shows only in the first. A long tail of unfinished pieces shows only in the second, and only as a large product with a small amount of real work behind it. A step that needs records from every other worker — a **wide step** — shows in the third. ## The ratio that actually diagnoses it The single most useful derived figure is **input bytes read divided by the bytes the result genuinely depends on**. A run that produces a daily summary of one day's records but reads three years of them has a ratio in the hundreds. Track its trend per job rather than its absolute level, because the absolute level is meaningless across different jobs and a step change in it is nearly always a real regression. | Symptom | Number that shows it | Likely cause | |---|---|---| | Huge bill, short run | input bytes read | the run read far more than its result needs | | Huge bill, long run, low work done | capacity held times duration | a few unfinished pieces holding the whole cluster | | Huge bill, moderate input | bytes moved between workers | records redistributed across the network at scale | | Bill grew with no code change | input bytes read, over time | the selection the job makes now matches more data | ## Where a tenfold gap usually hides - **A bound that quietly became open-ended.** A selection that once matched one day now matches everything ever stored, and the output is unchanged because the extra rows contribute nothing. - **The same input read more than once in one run.** Each pass is metered again. The mechanism for retaining an intermediate result instead of recomputing it belongs to the memory side of this subject; the bill is simply what makes the second pass visible. - **A tail of a few unfinished pieces.** Four hundred pieces finish in seconds and three run for an hour, and the run is billed for everything it still holds during that hour. - **An input that grew upstream.** The job did not change; the data did. This is the most common answer when two runs of the same code diverge. What you do about the first is largely a storage-layout and selection question owned elsewhere. What this leaf owns is the measurement that tells you it is happening at all. ## What varies between systems Do not assume every engine hands you these four numbers in the same shape. Some publish a run-level total for bytes read; others report only per-step input figures you must sum yourself, and a step that reads the same input twice may be counted once or twice depending on the engine. A continuous job built from repeated small finite runs reports per slice, so a figure for a billing period has to be aggregated across slices. Whatever the shape, capture the numbers as values that outlive the run — a number every worker adds to and the coordinating process sums, so the evidence survives the worker that produced it — rather than reading them from a live screen that will be empty tomorrow. ## How to use it Put one ratio on a recurring review per job — bytes read per unit of useful output — and alarm on a step change rather than on a threshold. A threshold has to be set per job and will be wrong; a doubling week over week is a signal in any job, and it usually arrives with a change you can name.

  • Why does a run reading far too much input pass every operating signal?
    Because those signals measure whether the job is working, not whether the work was necessary. Throughput is high, distance behind the input is zero, nothing failed or retried. Reading too much is correct behaviour on a wasteful instruction, so no health signal has a reason to complain. It only appears when bytes read are compared with the bytes the result depends on.
  • Two runs of the same unchanged code differ tenfold in bytes read. What explains that?
    The data the code selects changed. Usually a bound in the job's own selection stopped narrowing anything, so it now matches all stored history, or an upstream dataset grew wider or deeper. The output can stay the same size while the input multiplies, which is exactly why output size is a bad proxy for what a run cost.
  • Which single figure would you put on a per-job cost review?
    Bytes read per unit of useful output, tracked as a trend for that job. Absolute levels are not comparable across jobs and any threshold you set will be wrong for most of them, but a step change within one job is nearly always a real regression and usually lands in the same week as the change that caused it.

saying these in an interview costs you the question

  • Explains a cost gap from run duration alone
  • Assumes a run that never failed cannot be wasteful
  • Compares monthly totals instead of per-run measurements
  • Believes output size predicts what a run cost
  • Treats reading everything as free because each byte is cheap
  • Sets an absolute threshold instead of watching the trend per job