A job keeps everything it remembers between records as live objects inside the worker process - what caps that, and what happens at the cap?
answer
- one process, not the cluster
- shared with the steps' working memory
- no overflow path to disk
- reclamation pauses, then the crash loop
basics
~20 sThe cap is the memory of that one worker process, shared with everything else the process is doing. There is normally no overflow path, so passing it degrades and then kills the worker rather than quietly moving entries to disk.
solid answer
~50 sEverything a running job still holds between records - running accumulators, records buffered awaiting a group, deduplication entries, join sides, pending per-key wake-ups - is its `retained set`. One place that set can physically sit is in-process memory: ordinary live objects inside the worker process, the process running some of the job's parallel steps, so a read is a pointer dereference with no encoding at all. The cap is that single process's memory, not the cluster's total, and the retained set shares it with the steps' own working memory. Crucially there is usually no overflow to disk from this placement: the process gets slower as automatic memory reclamation runs harder, then exceeds its limit and dies, and the restarted worker meets the same data again. Some engines offer a disk-backed placement as the alternative; some fix the placement for you.
go deeper
Recall that state lives inside one worker process, and that its limit is that process's memory rather than the cluster's. Know that a filter holds nothing while a running total per key holds something for every key.
Explain why there is no overflow path from this placement, and why latency rises from memory reclamation before the process dies. Be able to contrast it with a disk-backed store in one sentence.
Show that you plan for the restart loop, not just the crash: describe how you would detect the approach of the ceiling and what you would change first, given that more memory only moves the line.
Frame it as a standing rule: which jobs in the estate are allowed to hold state in process memory at all, what retained size forces a different placement, and who owns the decision before the first deployment.
## What is being held, and why it has to live somewhere A long-running job remembers things between records. Call all of it **the retained set**: running accumulators, records buffered until the group they belong to is released, deduplication entries, both sides of a join, and pending per-key wake-ups the job registered for later. A filter or a projection remembers nothing and has no retained set at all; a running count per customer remembers one value per customer and keeps remembering it until something removes it. That set is not an abstraction - it is bytes on a machine. One of the places it can sit is **in-process memory**: the values are kept as ordinary live objects inside **the worker**, the process running some of the job's parallel steps and holding their state. A read is then a pointer dereference. Nothing is encoded, nothing is decoded, nothing touches a disk. It is the fastest placement anyone offers, and it has the hardest ceiling. ## Where the ceiling comes from - The ceiling is **one worker process's own memory**, not the cluster's. Ten workers with 8 GB each do not give one worker's keys 80 GB; they give each key whatever the single process that owns it has. - The retained set does not get that memory to itself. The same process is running the job's steps and holding their working memory too, so the share available to stored entries is smaller than the process limit. (How that division is made is a separate subject.) - On runtimes that reclaim memory automatically, a large **live** set costs more than the same bytes sitting idle: the reclaimer repeatedly walks objects that are all still reachable, and pause time grows with the live set rather than with allocation rate alone. - Local disk on the worker does **not** extend this placement. Writing to that disk is what the other placement does deliberately; the in-process memory placement has no such path. ## What failure actually looks like 1. **Latency degrades first.** Reclamation runs more often and returns less each time, so per-record processing time rises while nothing about the input changed. 2. **The process crosses its limit and dies.** Depending on the deployment, the runtime raises an out-of-memory failure or an external limit terminates the process. 3. **The restart repeats it.** The job restores the retained set from its durable snapshot - the periodic copy written outside the workers - reaches the same data, and hits the same ceiling. The result is a crash loop, not a single failure. The important property is the absence of a graceful middle. A disk-backed placement gets slower as it grows; this one runs at full speed until it does not run at all. ## The placements side by side | placement | what a read costs | ceiling | behaviour at the ceiling | |---|---|---|---| | live objects in the worker process | a pointer dereference | that one process's memory, shared with its steps | degradation, then the worker dies | | embedded on-disk store with a memory cache | an encode and a decode, plus a disk read when uncached | the worker's local disk | gradual slowdown as the cache covers less | | changed entries written as versioned files each cycle | paid at cycle boundaries, not per access | shared storage, effectively far larger | longer cycles, more files to manage | ## What raises the ceiling, and what only looks like it does - **More memory per worker** genuinely raises it, but does not change the growth rate, and a very large live set makes reclamation pauses worse. - **More workers** splits the keys across more processes, so each holds less - but only to the extent the keys spread evenly, and only in whole buckets, because a key's bucket is fixed when the job first starts. - **A bigger local disk** does nothing for this placement. It only matters once you move to a disk-backed one. - **A faster network or a bigger cluster** does nothing either; the constraint is one process's memory. ## What varies between engines Do not assume this choice exists everywhere. Some engines expose both an in-process memory placement and a disk-backed one and let the author pick; some expose only one. The oldest model in this family - the two-phase disk-to-disk batch model that writes intermediate results to disk between two fixed phases and keeps nothing between runs - has no long-lived retained set, so the question does not arise for it at all. A runtime that cuts a continuous input into short finite chunks and runs a complete job over each one typically keeps state as versioned files rather than as live objects, which shifts the cost from a per-process ceiling to a per-cycle write. When you answer this in an interview, say which model you are describing; a claim that is exactly right for one of them is simply false for the next.
- Does giving the worker more memory solve it?It moves the ceiling, not the growth rate, so it buys time rather than a fix. A much larger live set also makes automatic memory reclamation pauses longer, and a bigger retained set lengthens both the durable snapshot and the restore that follows a restart. Treat extra memory as headroom while a removal rule or a different placement is put in place.
- Why is the failure abrupt rather than a gradual slowdown?Because this placement has no second tier. A disk-backed store degrades smoothly as more accesses miss its cache, but live objects either fit in the process or do not; past the limit the process is terminated. The only warning is rising latency from memory reclamation shortly before the end.
- Does adding workers always reduce what each one holds?Only when the keys spread across them, and only in whole buckets: a key's bucket is fixed at the job's first start, so a parallelism change moves buckets between workers rather than re-splitting keys. If most of the retained bytes sit under a few keys, the worker owning them holds nearly as much as before.
saying these in an interview costs you the question
- Thinks the ceiling is the cluster's total memory rather than one worker's
- Expects entries to move to local disk automatically when memory runs short
- Believes only the part touched each cycle has to fit in memory here
- Assumes the retained set is the only user of the worker's memory
- Says doubling the worker count always halves each worker's retained bytes
- Treats the crash as a one-off rather than a repeating restart loop