skip to content

Memory or Local Disk

Holding state in the worker's own memory against an embedded on-disk store with a cache: the ceiling of one, the read and compaction cost of the other, and how to choose between them.

on this pageshow

questions

4

A job keeps everything it remembers between records as live objects inside the worker process - what caps that, and what happens at the cap?

level: juniorimportance: must knowfreq 58%

answer

  1. one process, not the cluster
  2. shared with the steps' working memory
  3. no overflow path to disk
  4. reclamation pauses, then the crash loop

basics

~20 s

The 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 s

Everything 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

for a junior

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.

for a middle

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.

for a senior

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.

for a principal

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
open as a page

What changes when a stateful step moves its retained set from process memory into an embedded on-disk store with a memory cache?

level: middleimportance: must knowfreq 66%

basics

~20 s

The ceiling moves from one worker process's memory to its local disk, and in exchange every read and write pays an encode and a decode, a disk read whenever the cache misses, and the background merging of the store's own files.

open as a page

A worker retains 40 GB of state but touches only 2 GB per cycle and has 16 GB of memory - which placement fits, and what could overturn it?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Forty gigabytes cannot live as objects in a sixteen-gigabyte process, so a disk-backed store with a memory cache fits, and the two gigabytes touched per cycle are mostly cache hits. Restore time or an unstable hot set can overturn it.

open as a page

A continuous job runs as repeated small finite jobs and writes changed state entries to a versioned file set each cycle - what dominates its state cost?

level: seniorimportance: nice to knowfreq 32%

basics

~10 s

Cost is paid per cycle, not per access: each cycle writes a file per state partition to durable shared storage, plus periodic full copies, so partitions multiplied by cycle rate sets the bill.

open as a page