skip to content

The Retained Set

The accumulators, buffers, deduplication records and join sides a continuous computation keeps between records, and why the retained set outgrows what the input rate suggests.

on this pageshow

questions

4

A continuous job runs a filter, a per-user running total and a repeat-identifier dropper — which of the three hold something between records?

level: juniorimportance: must knowfreq 70%

answer

  1. depends on more than this record
  2. filter forgets, total remembers
  3. seen-before is a question about the past
  4. one entry per distinct key
  5. buffers, join sides, pending wake-ups too

basics

~20 s

The running total and the repeat-identifier dropper hold entries between records: one accumulator per user, one entry per identifier already seen. The filter holds nothing, because its verdict depends only on the record in hand.

solid answer

~50 s

The test is whether the output for the record in front of the step depends on anything other than that record. The filter's verdict comes from one field of one record, so it keeps nothing between records. The running total needs the previous total for that user, which is not in the record, so it holds one accumulator per distinct user. The dropper has to answer "have I seen this identifier before", which is a question about the past by definition, so it holds one entry per distinct identifier admitted so far. Everything a job is still holding between records is what this tree calls **the retained set** — running accumulators, records buffered awaiting a group, deduplication entries, both sides of a join, and pending per-key wake-ups. The first two steps contribute to it; the filter does not.

go deeper

for a junior

Recall the test: if the output needs anything beyond the record in hand, the step is holding something. Name the obvious holders — a running total and a seen-before set — and say a filter holds nothing.

for a middle

Explain the full inventory, not just accumulators: buffered records awaiting a group, both sides of a join, and pending per-key wake-ups are all held, and each grows on a different count.

for a senior

Show that you classify steps at design time to predict cost, and say when the retained set is touched under the model you are describing — once per record, or once per finite chunk of an endless input.

for a principal

Frame the classification as the input to a size and restart-time budget for a fleet of jobs, and insist that a design says which processing model it assumes rather than leaving it implied.

## The one test A step in a long-running job is **stateful** if the answer it produces for the record in front of it depends on something other than that record. If the answer can be produced from the record alone, plus constants fixed when the job was built, the step is **stateless** and keeps nothing between records. Everything a running job is still holding between records is called here **the retained set**: running accumulators, records buffered awaiting a group, deduplication entries, both sides of a join, and pending per-key scheduled callbacks — wake-ups the job registered against one key and one moment, which are themselves stored like any other entry. Apply the test to the three steps in the question: - **The filter** — "drop records whose amount is zero". It reads one field of one record and decides. The previous ten million records cannot change the verdict. It contributes nothing. - **The per-user running total** — the answer is the previous total for that user plus this amount. The previous total is not in the record, so it must have been kept: one accumulator per distinct user. - **The repeat-identifier dropper** — "emit this only if I have not already seen this identifier". The set of identifiers already admitted is exactly the thing being held: one entry per distinct identifier. ## What the retained set contains - **Running accumulators** — a sum, a count, a maximum, a small sketch, one per grouping key. - **Records buffered awaiting their group** — where a computation cannot fold a record into a single value and must keep the records themselves until the group is released. - **Deduplication entries** — one entry per identifier already seen, held so a repeat can be recognised. - **Join sides** — records from each input held until a record from the other input can be matched against them. - **Pending per-key scheduled callbacks** — a wake-up registered against one key and one moment; it occupies storage from the moment it is registered until it fires or is cancelled. ## What holds nothing, and the one thing that looks stateless but is held Filters, projections, field renames, format conversions and arithmetic over one record hold nothing between records. One case looks stateless and is not: a step that enriches each record from a small lookup copy given to every worker is holding that copy — **per-worker state**, attached to one parallel instance of a step rather than to a key. It is held, but its size is set by the lookup, not by the input, so it does not grow as the job runs. | Step | Answer depends on | Held between records | Count grows with | |---|---|---|---| | Filter on one field | that record only | nothing | nothing | | Per-user running total | that record and the user's previous total | one accumulator per user | distinct users | | Repeat-identifier dropper | that record and every identifier seen before | one entry per identifier | distinct identifiers | | Enrich from a shipped lookup | that record and a fixed table | one copy of the table per worker | the table, not the input | ## Where the engines genuinely differ The classification above holds across this whole class of engine, but *when* the retained set is touched does not: 1. Runtimes that handle **each record as it arrives** read and write the retained set once per record, so per-access cost is paid a billion times. 2. Runtimes that cut an endless input into **short finite chunks and run a complete job over each** load and write the retained set once per chunk, so the same logical entries are touched a few thousand times instead. 3. The **two-phase disk-to-disk batch model** keeps nothing at all between runs: a per-user total there is recomputed from the whole input every time, and the question of what is retained between records simply does not arise. Engines also differ in what they demand before a step may hold an entry under a key. Some require the stream to be grouped by that key first, so the entry is **key-bound state** — readable and writable only under the grouping key of the record being handled, held by the worker that owns that key, with no coordination needed for a read. Others let a step accumulate per parallel instance instead. A candidate should state which model they are describing rather than assume it. ## Why an interviewer opens here Because the classification predicts cost. A stateless step's memory is a function of how many records are in flight; a stateful step's is a function of how many distinct keys are live and how long their entries stay. The retained set is also the part of the job that has to survive a restart, so it is what a **durable snapshot** — the periodic consistent copy written to storage outside the workers — has to carry. Getting the inventory wrong at design time is how a job that passes its test dies in its third week.

  • The dropper only ever sees each identifier a few times. Does that make its retained set small?
    No. What it holds is one entry per **distinct** identifier admitted, not per occurrence. If identifiers are near-unique, the entry count is essentially the number of records admitted so far, which is the largest retained set of the three steps by a wide margin.
  • Is a step that assigns each record a new field from a fixed table stateful?
    It holds the table copy, but it is not stateful in the sense that matters for sizing: the copy is bounded by the table and does not grow as records arrive. Sizing questions care about entries whose count tracks distinct keys or records, not about a fixed payload given to each worker.
  • Does a job with no stateful step still have to hold anything?
    Usually a read position per input piece, so it can resume where it stopped, and whatever records are in flight. That is per-worker bookkeeping, bounded by the number of input pieces rather than by the data, and in the two-phase disk-to-disk batch model even that does not persist between runs.

saying these in an interview costs you the question

  • Says any step that reads a field is stateful
  • Thinks a filter holds the records it dropped
  • Counts only accumulators and forgets deduplication entries
  • Believes a pending per-key wake-up costs no storage
  • Assumes every engine keeps something between records
  • Calls a lookup copy given to each worker stateless
open as a page

A continuous job ingests only two thousand records a second yet retains a few hundred gigabytes — which quantities actually set that size?

level: middleimportance: must knowfreq 62%

basics

~20 s

Distinct live keys, bytes per entry and per-entry overhead set the size, with the retention horizon deciding which keys count as live. Input rate never enters for accumulators or deduplication entries, so a slow job can retain hundreds of gigabytes.

open as a page

A team sized its retained set by counting per-key accumulators alone and came out tenfold low — which held entries did that count miss?

level: seniorimportance: should knowfreq 48%

basics

~20 s

It missed records buffered awaiting their group, both sides of a join held until a match is possible, deduplication entries, pending per-key wake-ups, and the per-entry overhead the holder adds. Each is counted on a different quantity from accumulators.

open as a page

Across dozens of long-running jobs, what standing rule would make each team predict its retained set before first deployment?

level: principalimportance: should knowfreq 38%

basics

~20 s

Require a written retained-set budget per job before first deployment: entry categories, live key cardinality with its source, the horizon per category, bytes per entry, and the product. Fix the units so numbers compare, and require re-declaration when the grouping key or horizon changes.

open as a page