What changes when a stateful step moves its retained set from process memory into an embedded on-disk store with a memory cache?
answer
- pointer read against encode and decode
- ceiling moves to the worker's local disk
- cache hit rate, not total bytes
- background file merge is the hidden cost
- local files die with the worker
basics
~20 sThe 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.
solid answer
~50 sAn **in-process memory store** keeps the retained set - everything the job holds between records - as live objects in the worker, so a read is a pointer dereference and the hard limit is that process's memory. An **embedded on-disk store with a memory cache** is a key-value store running inside the same worker that keeps recent entries in memory and the rest in local files it merges in the background. Three things change. The ceiling becomes the worker's local disk, so the retained set may exceed memory. Every access now encodes on write and decodes on read, and misses the cache into a disk read, so performance tracks the cache hit rate rather than the total size. And the store does periodic background work merging its own files, which shows up as bursts of disk activity and occasional latency spikes with no change in the input.
go deeper
Recall the two places a retained set can sit on a worker: as live objects in the process, or in a key-value store inside the worker that keeps recent entries in memory and the rest on local disk.
Explain the exchange precisely: a higher ceiling bought with an encode and a decode per access, a disk read on a cache miss, and background merging of the store's own files. Say why the cache hit rate predicts performance better than the total size.
Demonstrate that you measure before choosing: access pattern and tail latency under a realistic key spread, plus what the placement does to snapshot duration and to how long a replacement worker takes before it serves its first record.
Own the default for the estate: which placement new jobs start on, what evidence is required to deviate, and how the fleet's recovery-time budget constrains how large any single worker's retained set is allowed to become.
## The two placements, stated by mechanism The retained set is everything a running job still holds between records: running accumulators, records buffered awaiting their group, deduplication entries, join sides, pending per-key wake-ups. It has to sit on a worker - the process running some of the job's parallel steps. - **In-process memory store.** The entries are ordinary live objects in the worker process. A read is a pointer dereference; nothing is encoded. The ceiling is that one process's memory, and there is no overflow path. - **Embedded on-disk store with a memory cache.** A key-value store runs inside the same worker process. Recent or frequently read entries sit in a memory cache; the rest sit in files on the worker's own local disk, which the store rewrites and merges in the background so that reads do not have to consult an ever-growing pile of them. Both are **worker-local**. Neither is the durable copy: the periodic snapshot written outside the workers is a third thing, and the bytes on the worker's local disk die with the worker. ## What one access costs | | in-process memory | embedded on-disk store with cache | |---|---|---| | read, hot entry | pointer dereference | cache lookup, then usually a decode | | read, cold entry | not applicable, everything is hot | disk read of one or more files, then a decode | | write | object mutation or replacement | encode, then an in-memory write that is later flushed to a file | | background cost | automatic memory reclamation over a large live set | merging the store's own files: extra disk reads and writes | | ceiling | the worker process's memory | the worker's local disk | The row that surprises people is the cache hit. A hit reliably saves the **disk read**; whether it also saves the **decode** varies, because many such stores cache encoded blocks rather than decoded values while others cache the values. Plan on the decode being paid, and treat anything better as a bonus. ## Why the cache hit rate, not the size, predicts performance A disk-backed store holding 40 GB behind a 4 GB cache can behave almost like memory or almost like a disk, and the size alone does not tell you which: - If each cycle touches a small, stable subset - the hot working set - nearly every access is a cache hit and the placement is cheap. - If accesses are scattered across the whole set, the hit rate approaches the cache-to-total ratio and nearly every access becomes a disk read. - Write-heavy patterns add a second effect: entries written once and read never still cost an encode, a file write, and their share of the background merge. This is why the honest answer to *is the disk-backed store fast enough* is always a measurement of the access pattern, never a number of gigabytes. ## The background file merge, and why it is not free The store writes new and updated entries into new files and marks old versions superseded. Left alone, a read would have to consult many files, so the store periodically merges them into fewer, larger ones. That work is real: 1. It reads and rewrites data that the job did not touch, so disk traffic exceeds what the job itself wrote. 2. It competes with the job's own accesses for the same disk and the same memory, producing latency spikes at the tail while the average looks fine. 3. It is usually asynchronous, which means it can fall behind; when it does, read cost rises because more files must be consulted. ## What it does to snapshots and to restarts The durable snapshot - the periodic consistent copy written to storage outside the workers, which a restarted job reads - is affected by the placement, though the mechanism of taking it belongs elsewhere. Where the store's files are immutable once written, a snapshot can ship only the files that are new since the last one, which makes snapshots of a large set far cheaper; a set held as live objects is generally serialised in full each time. Restart runs the other way: a disk-backed worker has to get files back onto local disk before it can serve reads, and while some engines warm lazily or in the background, the state has to be materialised eventually and that time counts against recovery. ## What varies between engines Not every engine offers both placements, and not every engine calls the trade the same way. A runtime that handles each record as it arrives pays the per-access cost a billion times; a runtime that cuts a continuous input into short finite chunks and runs a complete job over each touches state once per chunk, so the same per-access cost matters far less to it and its state may not sit in a worker-local store at all. The two-phase disk-to-disk batch model keeps nothing between runs, so neither placement applies. State the model you mean before you state the trade-off.
- Does a cache hit make the disk-backed store as fast as live objects?No. A hit reliably removes the disk read, but most such stores cache encoded bytes, so the decode is still paid on every read and the encode on every write. Assume a hit is cheaper than a miss by a large factor and still more expensive than a pointer dereference; where a store caches decoded values the gap narrows, so say which you are assuming.
- The worker holding the state store dies. What is on that local disk worth?Nothing on its own. Those files are a working copy that lives and dies with the worker; the copy that matters is the durable snapshot written to storage outside the workers. A replacement worker reconstructs its local files from that snapshot, which is why snapshot and restore cost belong in the placement decision.
- Which workload shape suffers most from the per-access cost?Record-at-a-time processing with several state reads and writes per record: the encode and decode are paid once per access, so the cost scales with record count. A runtime that processes a continuous input as repeated short finite jobs touches the state once per chunk, so the same per-access cost is amortised over many records and rarely dominates.
- When is the in-process memory placement still the right answer?When the whole retained set fits comfortably in the worker process with headroom for growth in distinct keys, and the latency target is tight enough that per-access encoding matters. Buy it deliberately, with a removal rule in place, knowing the failure at the ceiling is a dead worker rather than a slowdown.
A bench where every tool is already in your hand, against a cabinet standing beside the bench. The cabinet holds far more than the bench ever could, but each tool comes out in a case you have to open and close, and every so often you stop work to reorganise the drawers so you can still find anything.
saying these in an interview costs you the question
- Says the disk-backed store removes the memory constraint entirely
- Believes a cache hit costs the same as a live-object reference
- Treats the worker's local state files as the durable copy
- Thinks background file merging is free because another thread does it
- Compares the placements on average latency and ignores the tail
- Assumes every engine in this class offers both placements