A continuous job ingests only two thousand records a second yet retains a few hundred gigabytes — which quantities actually set that size?
answer
- count entries, not records per second
- distinct keys, not throughput
- horizon decides which keys are live
- per-entry overhead at a billion entries
- rate enters only as records per key
basics
~20 sDistinct 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.
solid answer
~50 sSize **the retained set** — everything the job still holds between records — from the entry count, not the throughput. For running accumulators and deduplication entries, the count is the number of **distinct** keys still live, where the retention horizon decides which keys count, and the size is that count times bytes per entry plus a per-entry overhead the holder adds for its own bookkeeping. Input rate does not appear in that product at all: forty million distinct accounts each holding a quarter-kilobyte accumulator is ten gigabytes whether the job sees two thousand records a second or two hundred thousand. Rate enters in exactly one place — for records buffered awaiting a group and for join sides, the count is records **per key** inside the horizon, so rate multiplied by horizon does set those. That is why a modest-throughput job with high key cardinality and a week-long deduplication horizon dwarfs a fast job over a handful of keys.
code
json · 25 lines{
"inputs": {
"recordsPerSecond": 2000,
"distinctLiveAccounts": 40000000,
"deduplicationHorizonDays": 7,
"joinMatchHorizonHours": 24
},
"heldCategories": [
{ "name": "accumulator per account", "entries": 40000000, "payloadBytes": 250, "payloadGb": 10 },
{ "name": "deduplication entry", "entries": 1210000000, "payloadBytes": 64, "payloadGb": 77 },
{ "name": "join side records held", "entries": 173000000, "payloadBytes": 500, "payloadGb": 86 },
{ "name": "pending per-key wake-up", "entries": 40000000, "payloadBytes": 50, "payloadGb": 2 }
],
"totals": {
"entries": 1463000000,
"payloadGb": 175,
"assumedOverheadBytesPerEntry": 40,
"overheadGb": 59,
"retainedGb": 234
},
"notes": [
"recordsPerSecond appears only inside horizon-derived entry counts",
"overheadBytesPerEntry is an assumption to be measured, not a published constant"
]
}go deeper
Recall that what a job holds is counted in entries, not in records per second, and that one entry per distinct key is the starting point for any estimate.
Produce the arithmetic out loud: live key count times payload plus per-entry overhead, with the retention horizon deciding which keys are live, and say where input rate does and does not enter.
Be explicit that total retained bytes is the number driving restart and snapshot duration, keep it separate from the working set touched per interval, and treat per-entry overhead as measurable rather than guessed.
Require the estimate before first deployment and require its assumptions — cardinality source, horizon, bytes per entry — to be written down, because the number is only as good as the cardinality claim underneath it.
## The arithmetic, stated precisely For each category of held entry: **total bytes = (number of live entries) x (payload bytes per entry + per-entry overhead)** Everything interesting is in the first factor, and it is different per category: | Held category | Entry count is driven by | Does input rate enter? | |---|---|---| | Running accumulator | distinct keys still live | no | | Deduplication entry | distinct identifiers admitted inside the horizon | only via how many distinct identifiers appear | | Records buffered awaiting a group | distinct keys x records per key inside the horizon | yes, as records per key | | Join side | distinct keys x records per key inside the match horizon | yes, as records per key | | Pending per-key scheduled callback | keys with a wake-up outstanding | no | So the headline — *the retained set is not a function of input rate* — is exactly right for accumulators and pending wake-ups, and is a dangerous simplification for buffers and join sides, where rate multiplied by the horizon is the count. Say which you mean. **Retention horizon** is what decides whether a key counts as live at all: a key touched once and never again is still held until something removes it, so a horizon of a week and a horizon of a year give wildly different live counts over the same input. What removes an entry, and on which clock, is another node's subject; here the horizon is an input to the arithmetic, not a policy question. ## Why the intuition fails The instinct is to reason from throughput because throughput is the number on the dashboard. But throughput sizes what is **in flight** — the records a worker is currently handling — which is bounded by parallelism and the size of a record, and is usually small. The retained set is sized by **accumulation**, and accumulation is a function of how many distinct things the job is tracking and for how long. Two jobs at identical throughput can differ by four orders of magnitude in what they hold: one keyed on a hundred merchant categories, one keyed on a near-unique request identifier. ## A worked estimate Take the job in the question: two thousand records a second, keyed on account, forty million distinct accounts active, deduplicating on a near-unique identifier over a seven-day horizon, and matching against a second input with a twenty-four-hour horizon. 1. **Accumulators** — 40 million live accounts x 250 payload bytes = about 10 gigabytes. Rate is nowhere in this line. 2. **Deduplication entries** — seven days at two thousand a second is about 1.21 billion near-unique identifiers; at 64 payload bytes that is about 77 gigabytes. Rate enters only because it determines how many distinct identifiers arrive. 3. **Join side** — twenty-four hours at two thousand a second is about 173 million records held; at 500 bytes that is about 86 gigabytes. 4. **Pending wake-ups** — one per live account at about 50 bytes is about 2 gigabytes. 5. **Per-entry overhead** — roughly 1.46 billion entries in total; at a few tens of bytes each for the holder's own bookkeeping that is another several tens of gigabytes. The payloads alone come to roughly 175 gigabytes, and overhead pushes the total past 200. Nobody looking at two thousand records a second predicted that, and nothing in the arithmetic required a fast job. ## Per-entry overhead is not a rounding error Every holder adds bookkeeping per entry: a key copy, a length, an index position, alignment. At a billion entries, a few tens of bytes each is tens of gigabytes — comparable to the payload when payloads are small, and *larger* than the payload for a deduplication set of short identifiers. The exact overhead depends on how the holder represents an entry, and engines in this class differ several-fold on that point, so estimate it as a range and confirm it by measurement rather than quoting a number. Where those entries physically sit, and what that choice costs per access, is a different node's subject. ## What the number is for Total retained bytes is the number that drives how long a **durable snapshot** — the periodic consistent copy written outside the workers — takes to write and to read back on restart, and it is a design-time claim someone has to own. Distinguish it from the **working set**, the subset actually touched in a given interval, which is a different quantity used for a different decision. A job can have a small working set and an enormous total, or the reverse, and conflating them produces an estimate that is confidently wrong in whichever direction the author was thinking.
- Key cardinality doubles but throughput is unchanged. What happens to the retained set?It roughly doubles for every per-key category: accumulators, pending wake-ups and the per-key part of buffers and join sides all scale with distinct keys. Throughput being flat is irrelevant, which is precisely why cardinality growth, not traffic growth, is the number to watch on a long-running job.
- Which single input to the estimate is most often wrong?Distinct live key cardinality. Teams quote distinct keys per day when the retention horizon is a month, or quote a count from a sampled dataset. Get it from the source system's own cardinality over the actual horizon, and state the horizon next to the number so the two cannot drift apart.
- Does the same arithmetic apply to a job that cuts its endless input into short finite chunks?The entry counts are the same, because they are a property of the data and the horizon, not of the runtime. What changes is when the cost is paid: that model loads and writes the retained set once per chunk, so the per-access cost is multiplied by chunks rather than by records. Total size is unchanged.
A self-storage facility's bill is the number of lockers you rent, times how long you keep them, times the size of each. How many vans arrive at the gate per hour tells you nothing about it — a slow trickle of customers who never give a locker back costs far more than a busy day of people who load and leave.
saying these in an interview costs you the question
- Sizes the retained set from records per second
- Assumes a low-throughput job holds little
- Counts distinct keys but ignores per-entry overhead
- Treats total retained bytes and the touched working set as one number
- Forgets that a key touched once stays live
- Says rate never matters, including for buffers and join sides