Why is a counter that every worker adds to and a coordinating process sums better evidence than a line printed on a worker?
answer
- one number, not one record
- workers add, one process sums
- outlives the machine that produced it
- retries make a total approximate
- cardinality is the failure mode
basics
~20 sA carried counter is summed by the one process that outlives the workers and is normally written out with the run's figures, so the evidence survives the machine that produced it. Its limit is that it gives a total, never a record.
solid answer
~50 sA **carried counter** is a number every worker adds to locally and reports to the **coordinating process** — the single process that hands out pieces and tracks which finished — which sums the partial values into one total for the run. Because the total lands in the process that outlives the workers, and usually in whatever record the engine keeps of the run, it is readable after the machines are gone; a printed line is not. The price is resolution. A counter says `412 records rejected for a missing field`, not which 412. It must be declared before the run, it has to stay low-cardinality — one name per condition, never one per key — and its exactness under retry varies: some engines fold in only the attempt that completed, others count every attempt, so read a total as a strong signal rather than an audited figure.
go deeper
Know that a worker can add to a named number that the coordinating process totals for the whole run, and that this number is still readable after the workers are gone.
Explain the mechanics: local tally, reporting at a unit boundary, summation in one process, constant size regardless of data volume, and why per-key names break it.
Demonstrate judgment about exactness. Say which engine rule governs retried attempts, and move any figure that must reconcile out of telemetry and into the job's own output.
Argue for a fixed vocabulary of counters across every pipeline the organisation runs, so an unfamiliar job is legible at three in the morning by someone who has never seen it.
## What a carried counter is Every worker process holds a small local tally and adds to it as it handles records. At a reporting point — normally when a **unit of work** (the smallest thing the engine hands a worker, and the smallest thing it retries by itself) completes, sometimes on a timer — the worker sends its partial value to the **coordinating process**, which adds it to the run's total. That is the whole mechanism, and everything useful about it follows from where the total ends up. It ends up in the one process that is still alive when workers come and go, and typically in whatever figures the engine keeps for the finished run. A printed line lives on a machine the cluster will take back; a counter's value has already moved. ## Why the number outlives the machine - The worker is transient; the coordinating process spans the run. - Aggregation is cheap and constant in size — one integer per name, not one entry per record — so it can be reported continuously without loading the network. - Because it is small, it can be emitted while the job is healthy, which means the evidence exists *before* the first failure rather than being gathered after it. - It is the only run evidence that costs roughly the same whether a run processes a thousand records or a trillion. ## What a counter cannot tell you - **Identity.** A total gives volume, not the offending records. Pairing a counter with a small bounded capture of examples is a separate, deliberate arrangement. - **Location.** It says a condition happened, not in which piece of the input, unless you deliberately carry a second counter cut by something coarse. - **Cardinality.** One counter per distinct grouping key turns a constant-size report into a map as large as the key space, held in the single process whose death ends the run. That is not evidence; it is a new failure mode. - **A clean zero.** A counter reading zero can mean the condition never occurred, or that the records went down a path the counter does not sit on, or that the step holding it was never executed in this run at all. ## Retry, and the honest reading of a total When a unit of work fails, the engine normally retries it inside the same run — invisibly, from outside. If the failed attempt had already added to a counter, the total can include work that was thrown away. Engines differ in how they handle it: | Approach | What a total then means | |---|---| | Only a completed attempt's additions are folded in | A total close to exact for units that finished at least once | | Every attempt's additions are folded in | A total that overstates whenever anything was retried | | Additions inside a step that is recomputed for recovery | Can be applied more than once, for the same reason | So state the reading with its qualifier: `about 412 rejections, on a run with two retried units` is honest; `exactly 412` is not, unless you know your engine's rule. When a figure must be exact, the count belongs in the job's output — counted as data, through the same machinery that makes the output correct — not in the telemetry channel. ## Designing the small set worth carrying A practical default is a handful of totals, chosen so that the difference between them is itself informative: 1. **Records read** at the source step. 2. **Records emitted** at the final step. 3. **Records rejected, one counter per reason** — a malformed value is a different fact from a failed lookup. 4. Optionally **bytes read**, when the bill or the bottleneck is likely to be input volume. Read as a set they answer the first question after any incident: did the records go missing, or were they never there? Read one at a time they answer almost nothing. Two habits make them trustworthy: name each counter after the condition rather than the code path, so the meaning survives a refactor; and put the same names on every pipeline, so an unfamiliar run is legible to whoever is awake.
- A rejection counter reads zero, yet rows are missing from the output. What are the readings?Three, in order of likelihood: the records were dropped somewhere the counter does not sit, such as a filter earlier than the instrumented step; the step holding the counter was never executed in this run; or they were never read at all, which a records-read counter at the source would settle. A zero is evidence about one path only.
- When should the count live in the job's output instead of a carried counter?Whenever the number is a result rather than telemetry — anything reported to a person, reconciled against a source system, or billed on. Output goes through the machinery that makes the job correct, so it survives retries with the same guarantees the data has. A carried counter is diagnostic, and under retry it is approximate.
- Why not carry one counter per grouping key so you can localise anything?Because the number of names then grows with the key space, and every one of them accumulates in the single coordinating process that also hands out work. You trade a constant-size report for unbounded memory in the process whose death ends the run — and for coarse localisation, a few counters cut by a low-cardinality attribute do the job.
saying these in an interview costs you the question
- Believing a carried total identifies which records were affected
- Adding a counter per grouping key and expecting the summing process to cope
- Quoting a total as exact when units of work were retried
- Assuming counters appear without being declared before the run
- Reading a zero counter as proof the condition never occurred