clinker-exec · pipeline/memory.rs · ← back to Memory Arbitration

Keeping memory
at bay

Clinker promises bounded memory for finite batch jobs. No single mechanism delivers that. A budget, a ledger, an arbitrator, spill files, producer pauses and a memory-aware scheduler share the job. This page takes each one apart and gives you a control to try it. The simulator in section 05 runs them all at once.

01

The budgetOne number, three thresholds, two ways to measure

An author sets one value, pipeline.memory.limit (default 512M; the --memory-limit flag overrides it). The arbitrator turns it into three watermarks:

  • Soft limit, 80%. Pressure starts here. The arbitrator pauses a producer or asks a stage to spill.
  • Resume watermark, 70%. Paused producers restart once memory drops below this. Authors can tune it with resume_threshold, which must fall inside (0, 0.80); anything else is E324.
  • Hard limit, 100%. A request on the walk that does not fit here reclaims first, and the run stops with one E310 report only when that frees nothing. Three refusals run no round: a request made off the walk (a Source reading ahead, a writer, a join's own matching thread), one request larger than the whole limit, and Cull's per-group decisions.

The 20% gap between soft and hard is the spike allowance. Spill reacts at batch boundaries, not on every allocation, so short overshoots have to fit somewhere.

Which signal decides what. Pause, resume and the order a reclaim pass spills in read the charged total the registered stages report about themselves, so memory the process holds for other reasons never pauses a reader. Process RSS still counts in three places: its peak trips the spill poll, which then elects a victim; the spill before a paused Source resumes aims at whichever of RSS and the charged total sits higher above the soft limit; and a backstop check in the hot loops reads its peak, so a run whose unaccounted memory passes the limit says so rather than overrunning silently.
memory.rs · soft_limit(), resume_limit(), should_spill(), should_abort(), charged_bytes(), current_pressure()

Pressure gauge

02

The ledgerWho holds what: pull-mode attribution, plus reserve-before-allocate admission

Layer A: the consumer registry

Each stage that holds memory (Source ingest channels, hash Aggregates, sort buffers, grace-hash partitions, IEJoin arrays, Reshape and Cull group buffers, window arenas, and every inter-stage node_buffers slot) registers a MemoryConsumer with the run's single MemoryArbitrator.

A consumer holds a ConsumerHandle, a few atomics: live bytes, peak bytes, a spill-requested flag, a paused flag, and an active flag. The arbitrator never pushes numbers in. It pulls current_usage() from each consumer on every poll, so a grace-hash join with partitions on disk reports only what is still in memory.

A registration lasts as long as the state it describes. A Source's consumer is removed when its stream ends (its reader's end of input, or the walk stopping), and a buffer slot's when its last reader finishes. The registry therefore lists only live memory. The rows a finished Source read stay charged in its name until the steps holding them drop or spill them; the ledger keeps them as a finished Source's (retired_source), and an E310 report shows them as memory not held by any one node.

Layer B: the admission ledger

Newer storage (CSV decoding, prepared output, owned record storage) is stricter. An AllocationLease reserves the full Layout before anything is allocated. Growing a buffer reserves the new block while the old one is still charged. If the reservation is refused, the old contents and their charge stay as they were. Moving a grant between owners never leaves a moment where it is uncharged.

The two layers answer different questions. Layer A gives estimates the policy uses to choose what to spill. Layer B gives exact grants, and the arbitrator deducts them before admitting anything else. Neither layer is a promise about whole-process RSS; parser buffers and allocator overhead sit outside both.
memory.rs · MemoryConsumer, ConsumerHandle · engine docs: memory-arbitration.md "Exact allocation admission"

Try it: a ReservedVec under a 256 KiB budget

Allocate a vector, then double it. During growth, the old and new blocks are both charged at once, and this overlap is what most often hits the budget.

0budget 256 KiB
granted0
live vec cap0
peak grant0
03

ArbitrationUnder pressure, choosing which stage gives memory back

Every consumer has two fixed traits. Its spill priority (lower spills first) says how cheaply it can spill. Its can_back_pressure flag says whether its producer can be paused instead. Only Sources can be paused.

The memory.backpressure setting picks a policy:

knobpolicy objectpicks
pause (default)BackPressurePreferred → Prioritythe first pausable consumer, or else the lowest priority, most reclaimable on ties
spillPrioritylowest priority, most reclaimable on ties; never pauses
bothBackPressurePreferred → LargestFirstthe first pausable consumer, or else whoever a spill would free the most from

Both policies rank on reclaimable_bytes — what a spill of that consumer would free now, not what it holds charged. A consumer that reports 0, such as the output staging grant or a window arena, is never a candidate, so charged-only state never shadows one that can actually give memory back.

One round has two steps:
① reconcile_backpressure handles pause and resume, reading the charged total. It skips any Source the walk thread is reading from right now (pausing that one would deadlock the run).
② poll_arbitration asks the policy for a victim. If the victim cannot be paused, it sets the victim's spill flag, and the stage spills at its next batch boundary. If the victim is pausable, step ② does nothing, because step ① already handled pausing.

So under pause, while any Source is registered, step ② picks that Source and nobody gets a spill request from the arbitrator. Spill still happens, through the stages' own thresholds, the slot admission check and the reclaim passes described in section 04.

memory.rs · build_policy(), Priority, LargestFirst, BackPressurePreferred, poll_arbitration(), reconcile_backpressure()

Arbitration lab

Tick consumers to register them and set the bytes each holds. Each column runs the real selection code for one policy. Registration order matters for first pausable.

04

Backpressure & spillThree ways memory gets pushed back

WallBounded channels. A Source's ingest thread stops when its crossbeam channel is full, and streaming handoffs work the same way. This happens in the transport layer, and the arbitrator isn't involved.
BrakeArbitrator pause. Above the soft limit — measured in charged bytes — the first pausable Source that is not active waits on a Condvar. Below the resume watermark it starts again. Between the two thresholds nothing changes (hysteresis).
ValveSpill. State is written to disk as length-prefixed postcard frames, LZ4-compressed when the per-schema compress setting chooses to. Cumulative disk use is checked against storage.spill.disk_cap_bytes (E320).
FuseReclaim, then refuse. A request on the walk that does not fit runs reclaim passes first — spilling walk-owned state, the requester last — and retries after each. It stops the run with one E310 report only when a pass freed nothing and a final pass freed nothing either. A request off the walk, one larger than the whole limit, and Cull's per-group decisions are refused at once, with no round. Probe loops also check every 10K output records.

Where spill actually comes from

  1. Arbitrator request. try_spill sets a flag and does no I/O. The stage reads the flag with take_spill_request() at its next boundary and spills there. Buffer slots are served in a sweep before the next node is dispatched.
  2. Self-spill at admission. When a node_buffers slot is published and should_spill() returns true, the producer writes those rows straight to a spill file.
  3. Stage thresholds. A hash Aggregate spills when its group count reaches about 60% of its budget, or its value heap passes 40%. Sort, grace-hash, Reshape and Cull have similar thresholds of their own.
  4. Spill before resuming. When the walk reaches a Source that is paused, it first calls spill_reclaimable(over), which runs one reclaim pass over the walk-owned victims a spill would free bytes from, in the policy's order. over is how far the larger of process RSS and the charged total (current_pressure()) sits above the soft limit. Only then does it resume the Source and mark it active, so the resumed Source doesn't push memory straight back over the soft limit.
  5. A reclaim pass another request started. A short request anywhere on the walk elects victims for its own pass, so state that is nobody's own threshold — a sibling slot, the document dead-letter state's held rows, an Output's per-document buckets — spills for a request made elsewhere.
Pause and spill read different signals. peak_rss only ever rises (fetch_max), so once the run has crossed the soft limit every later spill poll counts as under pressure; spilling too often is safe. Pause and resume read the charged total instead, which can fall, so a reader is never left parked and memory the process holds for other reasons never pauses one.

Spill priorities (engine table)

consumerpriopausable
node_buffers slot0no
output staging (writer)0no charged-only
grace-hash Combine10no
Reshape · Cull15no
sort buffer · IEJoin build20no
sort-merge Combine25no
hash Aggregate · inline-hash Combine30no
Source ingestN/Ayes pause
streaming AggregateN/Ano
credential registry · scan materialization · window arenalastno last resort

The IEJoin sort buffer, the sort-merge Combine and the inline-hash Combine report nothing a spill would free now, so neither a poll nor a reclaim pass elects them: the two sort kernels spill on thresholds of their own, and the inline-hash table never spills. Their rows give the order they take once a pass can reach them.

Order follows reload cost. Re-reading a buffer slot is one postcard round-trip. A grace partition costs a little more. Reshape and Cull must re-split groups when they reload. A sort needs an external merge. A hash Aggregate costs the most to bring back.

engine docs: memory-arbitration.md "Per-operator arbitration parameters"
05

The whole machineA toy pipeline running the real decision rules

Two Sources start ingesting together. The walk thread drains orders into a buffer slot, aggregates it, then drains refs into a grace-hash build, then probes and writes. Each tick runs the spill poll with its pause/resume reconcile, flag servicing, the walk's requests (which reclaim before they refuse), the join's own limit check once it builds its tables and as it probes, and the Source reads, which are refused at once when they do not fit, following the code described above. Rates, sizes and the RSS overhead model are made up for teaching, so read the shapes of the curves, not the numbers.

Controls

backpressure
bytes on spill disk (compressed)disk_cap_bytes▲ spill · ◆ pause · ● resume markers on the memory chart

Arbitrator log

Things to try. On Default, nothing is paused: the RSS peak trips the spill poll first, the slot spills at admission, and the charged total never reaches the soft limit, which is the only reading that pauses. On Paused reader resumes, a large Aggregate pushes the charged total past the soft limit, so refs is paused while the walk is aggregating, after orders has been read. When its turn arrives, one reclaim pass spills first, aimed at how far the larger of process memory and the charged total sits above the soft limit, and then it resumes. Switch to Spill policy: nothing is ever paused, and the arbitrator now elects buffer slots and the build directly. Both stays close to Tight budget here, because it only differs from pause once no Source is registered; from then on it picks the biggest holder whatever its priority. Section 03 shows the difference directly. Oversized group shows the one shape spill can't fix: its request runs a reclaim pass and then a final pass, both free nothing, and only then is it refused. Join check reclaims shows where the grace-hash join checks the hard limit. It reads refs with no limit check, spilling its largest partition when the spill poll trips. Once refs is read, it builds a table for each partition still in memory, and the tables' index does not fit. Its check runs a reclaim pass, which spills the Aggregate, the index then fits, and the run completes.
06

The schedulerChoosing what runs next, before pressure builds

The walk runs one node to completion before starting the next. When several nodes are runnable together (independent chains feeding a later Combine), next_runnable picks one by minimising this tuple:

( !fits, −predicted_freed, −subtree_reclaim, stable_index )

  1. Headroom fit. Prefer a node whose predicted_peak fits within soft_limit − charged. An unknown peak (0B) always counts as fitting.
  2. Frees the most immediately. A blocking operator that is ready to drain beats a fresh Source.
  3. Frees the most eventually. Between fresh Sources, start the chain whose downstream Aggregate releases the most.
  4. Stable index. Topological order, which keeps the choice deterministic on every machine.
Predictions come from the plan and the on-disk sizes of input files (shown by --explain). Without them every term is 0, so the order is plain topological. Scheduling never changes output; it only affects how high memory peaks.
memory.rs · next_runnable() · engine docs: memory-arbitration.md "Scheduling"

Runnable frontier

07

When it failsThe diagnostics, and the limits of the guarantee

E310

Memory limit reached

A request did not fit. On the walk the engine first ran reclaim passes, and they freed nothing; off the walk, for one request larger than the whole limit, and for Cull's per-group decisions, no round ran. Common causes: one indivisible unit (a Reshape or Cull group, one aggregate row, a range-join block pair) is larger than the budget, or what fills the limit cannot be written to disk. The report names the node that asked and what the memory was for, the charged total beside private memory, the five largest holders with the state each is in, what the reclaim round asked and freed, a pasteable limit floor, and one remedy.

E312

Unsatisfiable budget

Under pause or both, a limit below the process's baseline RSS is rejected at startup. Otherwise the run would pause forever. spill skips this check and spills hard instead.

E320

SpillCapExceeded

Cumulative spill passed storage.spill.disk_cap_bytes. That's your own limit; the volume may still have space. The file that went over the cap is deleted on abort.

E321

Spill volume full

The disk itself ran out. Also: tmpfs /tmp is RAM, so spilling there frees nothing. Point storage.spill.dir at a real disk.

E324

Bad resume_threshold

The value has to fall strictly between 0 and the 0.80 soft limit, otherwise the hysteresis band doesn't exist.

Honest limits