A job deduplicates events by holding every id it has seen. Why does that retained set grow forever, and what bounds it?
answer
- insert-only logic
- nothing deletes an entry
- distinct ids, not record rate
- removal must be asked for
- expiry, callback, completeness
basics
~20 sA deduplication check only ever adds: each new id creates an entry and nothing deletes one, so the retained set grows with distinct ids for the life of the job. Only a removal rule bounds it.
solid answer
~50 sThe retained set is everything a running job is still holding between records — running accumulators, records buffered awaiting a group, deduplication entries, both sides of a join, and pending per-key callbacks. A deduplication check is pure insert: look the id up, and if it is absent, store it and emit the record. There is no step in that logic that removes anything, and no runtime removes state you did not ask it to remove, so the entry count tracks the distinct ids the job has ever seen. Over an endless input that number never stops rising. Bounding it means adding a removal rule: entry expiry, which drops an entry a fixed interval after its last write or last read; a per-key scheduled callback that clears one key at a moment the job registers; or cleanup once the job can assert no older record is still expected.
go deeper
Recall that a deduplication check stores every id it sees and never deletes one, so the set grows with distinct ids for the life of the job. Know that something has to remove entries and that the author must ask for it.
Explain the sizing: distinct keys times bytes per entry, with the input rate absent from the formula. Name the three removal rules and say which clock each is measured against.
Show that you plan the removal rule before the first deployment, and that you know growth surfaces first as a restart that will not finish, not as an out-of-memory failure.
Frame it as a standing convention: every stateful step in the estate declares its key space and its removal rule at review time, because the cost of discovering a missing one is paid months later by whoever is on call.
## The retained set, and which parts of it can grow **The retained set** is everything a running job is still holding between records: running accumulators, records buffered awaiting a group, deduplication entries, both sides of a join, and pending per-key callbacks. Treat it as the job's real storage system — it has a size, a per-access read and write cost, a durability story, and, only if the author gives it one, a removal policy. Not every step retains anything. A filter or a projection looks at the record in front of it and forgets it, and can run for a decade at a constant footprint. A **deduplication check cannot**: to know whether it has seen an id before, it must still be holding that id. ## Why the logic has no natural end Written out, the check is pure insert: 1. Take the id off the record. 2. Look it up among the entries this step holds. 3. If it is present, drop the record as a duplicate. 4. If it is absent, store it and emit the record. There is no fifth step. Nothing there deletes, and no runtime in this class deletes an entry you never told it to delete. So the entry count tracks **the number of distinct ids the job has ever seen**. Over a finite input that number stops at the end of the input; over an endless one it does not. The sizing arithmetic follows, and this is the part candidates most often get wrong: - **entries** ≈ distinct ids within whatever retention you actually enforce — with no removal rule, that is all of them, since the job began; - **bytes** ≈ entries × (encoded key + encoded value + whatever per-entry overhead the store adds, which is rarely zero); - **the input rate does not appear in either line**. A million records a second over a thousand repeating ids is a thousand entries. Ten records a second of globally unique ids is ten new entries a second, forever. ## The three rules that actually remove an entry | removal rule | what triggers it | where it fits | |---|---|---| | **entry expiry** | a fixed interval has passed since the entry's last write, or in the other variant since its last read | deduplication sets, lookup-style entries, anything with a natural staleness horizon | | **per-key scheduled callback** | a moment the job registered against one key arrives, and the handler deletes that key's entry | a burst of activity from one user ended by a gap of inactivity; anything with a per-key deadline | | **cleanup on a completeness claim** | the job asserts no record older than a given moment is still expected, so entries belonging before it can go | grouped results that can be declared closed; the clock that produces the assertion is another node's subject | All three are things the author asks for. None of them is on by default in a way you can rely on. ## What differs between engines, and why it matters here This is a class of systems that disagrees about mechanism, so state the model you are describing rather than assuming one: - **Which rules are offered differs.** Some runtimes express removal as a declarative expiry attached to the entry; others give you only a callback you register and a handler you write; several offer both, and the two behave differently on restart. - **The clock differs.** An interval may be measured against the worker's own wall clock, or against the job's own assertion of how far the input has progressed. The second stops advancing when the input does, and entries then stop being removed. - **Removal may be lazy.** Some stores drop an expired entry the moment it is due; others mark it and reclaim the space only when the entry is next touched or during a background file merge, so bytes fall later than the entry count does. - **Some models have no long-lived retained set at all.** The two-phase disk-to-disk batch model writes intermediate results to disk between its two phases and keeps nothing between runs, so unbounded growth of this kind cannot happen in it. A runtime that cuts a continuous input into repeated small finite jobs touches its state once per chunk rather than once per record, which changes the cost of removal but not the need for it. ## It is not only a memory problem The first instinct is that the worker will eventually run out of room. That happens, but it is usually not the first symptom, and more memory only moves the date. Every byte in the retained set is also a byte in the periodic **durable snapshot** — the consistent copy written to storage outside the workers — and therefore a byte a restarting job must read back before it processes its first record. A job that has been fine for four months frequently announces the problem as a restart that no longer finishes. ## Saying it in an interview The answer an interviewer is listening for is two sentences of mechanism and one of sizing: the logic inserts and never deletes, so the retained set grows with distinct keys rather than with throughput; a removal rule has to be chosen deliberately, and choosing one means accepting the correctness it costs.
- Does storing a short hash of each id instead of the full id solve the growth?No. It lowers the bytes per entry by a constant factor, so the job reaches the same wall later on the same curve; the entry count is unchanged. It also introduces a collision chance, which here means silently dropping a genuine record as a duplicate. Fewer bytes per entry is worth having, but only a removal rule bounds the set.
- The source only keeps records for seven days. Does the retained set inherit that bound?No. How long the source keeps a record available governs what can still be read from it, not what the job is holding. An entry written by a record consumed a year ago is still in the retained set today unless a removal rule dropped it. The two horizons are independent, and conflating them is a common way jobs grow unnoticed.
- Which parts of a typical retained set can grow without bound, and which cannot?Anything keyed by an open-ended identifier can: deduplication entries, per-user or per-device accumulators, buffered join sides and pending per-key callbacks. Anything keyed by a closed set cannot — a per-country counter has at most a couple of hundred entries whatever the input does. The question to ask of every stateful step is what the key space is and whether it is finite.
A doorman writing every guest's name in a book and never crossing one out. The book's thickness tracks how many different people ever came, not how busy tonight is, and no amount of a bigger desk changes that — only a rule for tearing out old pages.
saying these in an interview costs you the question
- Thinks the worker evicts old entries automatically once memory is tight.
- Sizes the retained set from records per second instead of distinct keys.
- Assumes a restart clears it, ignoring that a snapshot restores every entry.
- Believes the source dropping old records also drops the job's entries.
- Treats it purely as memory, ignoring snapshot and restart time.
- Says a smaller per-entry encoding fixes growth rather than delaying it.