skip to content

A worker process is given one fixed memory budget. What competing demands share it while the job runs?

level: juniorimportance: must knowfreq 70%

answer

  1. one process, one wall, no lending
  2. several units of work share it
  3. borrowed for a step, then returned
  4. kept results are a second demand
  5. the runtime's own bytes go uncounted

basics

~20 s

Operators borrowing room while they run, any results the job asked the engine to keep, and runtime overhead the engine never counts - all shared by every unit of work inside that one process, with no lending from idle workers.

solid answer

~50 s

A worker is one operating-system process on one machine that runs some of the job's work and owns a fixed amount of memory nothing else can borrow. Three demands press on it. First, **operator working memory**: the room operators borrow while they run - sorting, building the in-memory table an aggregation fills, holding the built side of a join - and give back when the step ends. Second, **retained-result memory**: on engines that let you pin a computed intermediate for later reuse, the room those kept results occupy until they are dropped. Third, **off-budget overhead**: the language runtime itself, native and network buffers, and whatever a user-supplied function allocates outside the engine, which the engine's accounting never counted but the machine still holds. Several units of work usually run at once inside the worker and share all of it.

go deeper

for a junior

Recall the shape: one worker is one process with a fixed amount of memory nothing else can lend it, several units of work inside it share that memory, and part of what the process uses is never counted by the engine.

for a middle

Explain the lifetime difference - room an operator borrows for one step and returns, against a kept result that stays until dropped - and name the three different numbers people call memory.

for a senior

Show that you reason about which of the three numbers moved before touching anything, and that you know engines draw the region lines differently rather than assuming the one layout you have used.

for a principal

The angle here is what the platform standardises: the shape of a worker everyone inherits, whether authors may pin results at all, and who owns the consequence when that default is wrong for one team's workload.

## What a worker is, and why its budget is a wall A **worker** is one operating-system process, on one machine, that runs some of the job's work and owns a fixed amount of memory nothing else can borrow. A cluster of forty workers is forty separate walls, not one pooled budget forty times the size. A worker that is short of room cannot draw on the thirty-nine sitting idle, and adding machines does not change what one process can hold. That single fact is why so many memory conversations end in *divide the data into more pieces* rather than *give every machine more memory*. Inside that process, several **units of work** - one worker's share of one step, scheduled and retried on its own - normally run at the same time on separate threads, and they share the one budget. So *the worker has 8 GB* and *the worker is sorting 8 GB* are both sayable and cannot both be comfortable. Whenever you state a budget, state in the same breath how many units of work are sharing it. ## The three demands on that budget | demand | what it holds | when it is given back | counted by the engine? | |---|---|---|---| | operator working memory | sort buffers, the table an aggregation builds, the built side of a join | when the step that borrowed it ends | yes | | retained-result memory | intermediates the job asked the engine to keep for later reuse | when the result is dropped, evicted, or the job ends | yes | | off-budget overhead | the language runtime itself, native and network buffers, user-function allocations | only when the process exits | no | - **Operator working memory** is the borrowing region. An aggregation builds a **grouping table**: an in-memory table with one entry per distinct key seen so far, which grows with the number of distinct keys rather than with the number of records. A sort holds a run of records before ordering it. A join holds one side while it probes with the other. All of this is transient: the step ends and the room returns. - **Retained-result memory** exists where the engine offers a pin - telling it to keep a computed intermediate so later steps read it instead of recomputing the branch that produced it. This region is not universal. Some engines carve it out of the same pool as operator work; the oldest two-phase disk-to-disk model in this class has almost no notion of a kept result at all, and its memory story is a fixed working buffer per unit of work. - **Off-budget overhead** is the region that surprises people, because the engine's own reports do not show it. It includes the runtime of whatever language the worker runs, buffers allocated outside the engine for network and file work, and anything a user-supplied function allocates for itself - a lookup table loaded per unit of work, a large parsing buffer, a native library's own allocations. ## Three numbers answer to the word *memory* 1. **The engine's accounted budget** - what the engine divides into regions and tracks. This is what its reports are about. 2. **The process's actual footprint** - what the machine sees the process holding, which also includes the off-budget overhead. 3. **The limit the platform enforces** - a ceiling applied to the whole process by the platform rather than by the engine, and enforced by killing the process outright rather than by making one allocation fail. A junior is not expected to reason about the gap between them. A junior *is* expected not to use one word for all three. ## Size is not one number either The belief that a 200 MB input needs roughly 200 MB of memory is the specific error this subject exists to correct. The same records have at least three sizes: compressed at rest, packed in whatever layout the engine holds them in, and held as **host-runtime objects** - ordinary objects of the worker's language, each carrying the runtime's own per-object bookkeeping and each reclaimed automatically by the runtime, a process that stops the worker's own work for as long as it takes. The multiple between the last two runs several-fold. Always name the form when you state a size. ## What follows in practice - A worker short of room is a *local* condition. Nothing on the other machines is available to it. - What one unit of work holds is set by how much data its slice carries and which operators must hold a whole working set, not by the size of the cluster. - When an operator cannot hold its working set, engines that support it write part of it out to disk attached to the worker and read it back to finish - **spilling**. That is designed behaviour, not a bug, and how it is organised belongs to its own subject. - The engine's report is a report about region one. The machine enforces on region two plus three. The short version a first screen wants: one process, one budget, several units of work sharing it, three demands on it, and one of those three the engine cannot see.

  • Why does adding machines to the cluster often not help a worker that is short of room?
    Memory is owned per process. A new machine adds new workers with their own budgets; it does not enlarge the one that is struggling. It helps only indirectly, by letting the input be cut into more pieces so each unit of work holds less at a time. If the data is not divided differently, the same unit runs on the same-sized wall.
  • Which of the three demands does the engine's own report cover?
    Only the regions it manages - the room operators borrow, and kept results where the engine supports them. The overhead outside that accounting is invisible to the engine yet fully visible to the machine, which is why a report showing plenty of room free is not a statement about the process's real footprint.
  • Two units of work run at once in one worker and both sort. What happens to the budget?
    They draw from the same budget, so each effectively works against a fraction of it, and the fraction is not fixed in advance on every engine - some hand out room on demand, some reserve a share per concurrent unit. An operator that cannot get enough either spills to disk attached to the worker or fails, depending on the operator.

A one-room workshop. The bench you clear to work on a piece is cleared again when the piece is done; the shelf of finished pieces someone told you to keep stays occupied until someone takes them away; and the walls, the heating pipes and your own toolbox take floor space that the usable-area figure on the floor plan never counted. The fire marshal's occupancy limit, however, counts the whole room - which is how a workshop can be over its limit while the floor plan says there is space left.

saying these in an interview costs you the question

  • Thinks a busy worker can borrow memory from an idle worker
  • Assumes one unit of work has the whole worker budget to itself
  • Treats the engine's accounted budget as the process's real footprint
  • Says a 200 MB input needs roughly 200 MB of memory
  • Believes every engine keeps a region for results held for reuse
  • Thinks writing part of a working set to local disk means something broke