skip to content

When One Machine Is Not Enough

Measuring what a dataset really costs in memory, running the work in pieces when it will not fit, and knowing which operations refuse to be split. Most people guess, and guess low.

on this pageshow

explore

questions

18

Each of three steps on a large table returns instantly, then one later call runs for four minutes — why?

level: juniorimportance: must knowfreq 60%

answer

  1. nothing ran on the first three
  2. steps recorded, not executed
  3. one call carries the whole chain
  4. time the trigger, not the step

basics

~20 s

Plan-then-run execution: the three steps only recorded a description of the work, so they returned immediately. The later call is the trigger, where the whole recorded plan finally executes and every byte of work is paid at once.

solid answer

~40 s

The three steps did not compute anything. Under **plan-then-run execution** the library records each step as a description — read this source, keep these rows, add this derived column — and hands back a handle to the recorded steps rather than data. Nothing has been read yet, and appending a node costs the same whether the source holds a thousand rows or a billion. The fourth call is **the trigger**: the point at which the recorded steps are handed to the engine, the input is finally opened, and an answer appears. All four minutes belong to the whole chain, not to that one call. The practical consequence is that timing an individual step tells you nothing about the data, and the profile you actually need is of the trigger.

go deeper

for a junior

Recall that some libraries record a step instead of running it, and that the work then happens at one later call. A step returning instantly on a huge input is the tell.

for a middle

Explain what the handle actually holds — a description of steps, not rows — and why appending to it costs the same at any input size. Be able to say where the wall-clock time really belongs.

for a senior

Show how you would profile a chain whose per-line timings are meaningless, and say out loud that each prefix you time re-reads the source, so the measurement itself is not free.

for a principal

Frame it as a team-visible property: deferred entry points change what a timing, a stack trace and a log line mean, so the conventions for measuring and reporting pipeline cost have to be set alongside the choice.

The timing you observed is the clearest symptom of an execution model, and naming that model is most of the answer. ## Two models behind the same-looking code Under **eager evaluation**, every step runs on the line it is written. "Keep the rows where the amount exceeds 100" opens the input, evaluates the condition, allocates an object holding the survivors, and hands it back. The next step begins from that object. Time is spent where it is written, so a profiler's per-line numbers mean what they look like they mean. Under **plan-then-run execution** — deferred evaluation, where the steps are recorded rather than executed as they are written — those same three lines touch no data at all. Each appends a node to a description the library holds: *from this source; keep the rows matching this condition; add this derived column*. What comes back is a handle to the recorded steps, not a table of values. Appending a node is arithmetic over a handful of small objects, so it returns in microseconds whether the source holds a thousand rows or a billion. That is why your first three calls looked free: they were free, because they did nothing to the data. The fourth call is **the trigger** — the point at which the recorded steps are handed to the engine and an answer is demanded. Only then is the input opened. The four minutes are the cost of the whole chain, collected on one line. ## Reading the timing correctly | | eager | plan-then-run | |---|---|---| | when a step's work happens | on its own line | at the trigger | | what a step returns | a materialised object | a handle to recorded steps | | where wall-clock time appears | spread across the lines | concentrated at the trigger | | what a per-step timing means | the step's real cost | the cost of appending a node | | where a failure is reported | at the step that was wrong | at the trigger | The practical consequence: timing individual steps in a deferred chain measures nothing about the data. If you need to know which part of the work is expensive, use whatever the design offers for inspecting or instrumenting the recorded steps, or trigger progressively longer prefixes of the chain and diff the times — accepting that each prefix re-reads the input, because the earlier trigger retained nothing. ## Why a design would wait at all Deferral is not delay for its own sake. Holding the steps until an answer is demanded means the engine sees the end of the chain before it reads the first byte, and can use that knowledge in ways a step running on its own line cannot: - columns the recorded steps never reference need never be read or decoded; - a row condition written three steps in can be applied while the input is being read, so fewer rows are ever materialised; - consecutive element-wise steps can be collapsed into a single traversal, allocating one output rather than one full-length buffer per operator. None of that needs a network or a second machine. It follows from knowing what is wanted before starting. ## What the wait costs Two prices come with the same property, and both are visible in the timing you just saw. 1. **The failure arrives at the trigger, not at the step that was wrong.** The step never ran where it was written, so an error is reported at the line that executed the plan. How much that hurts depends on the design: some planners resolve the schema while the plan is being built and reject an unknown column immediately, deferring the work without deferring the check. 2. **A plan is a description, not a result.** Triggering the same handle again replays the chain from the source unless something was explicitly retained, so two answers drawn from one chain can cost two full runs. ## What varies between designs Not every library that defers does the same amount with the delay. Some record the steps and then execute them close to as written; others rewrite the chain substantially before running it. Some resolve column names and types at build time; others discover them at the trigger. Some offer an immediate and a deferred entry point over the same data model, so the identical expression has both timing profiles depending on which one you opened with. The answer an interviewer wants names the model first — *these three returned instantly because they were recorded, not run* — and then says which payoffs and which prices apply on the design in front of you, rather than asserting one tool's behaviour as the rule. The diagnostic itself is simple. If three heavy-looking steps return instantly and one later call carries all the time, you are in a plan-then-run model. If the time spreads across the lines roughly in proportion to the work each one describes, you are in an eager one.

  • How would you find out which part of a deferred chain is actually expensive?
    Use whatever the design exposes for inspecting or instrumenting the recorded steps first, since that costs nothing. Failing that, trigger progressively longer prefixes of the chain and diff the elapsed times — remembering that each prefix re-reads the input from the source, so the measurements are cumulative rather than independent, and the total cost of the exercise is several full runs.
  • If the steps returned instantly, does that mean the deferred version is doing less total work?
    Not by itself. Instant returns only mean nothing had run yet. The chain may genuinely do less work — unreferenced columns never read, fewer rows materialised, fewer intermediate buffers — but that comes from what the engine does with the whole chain at the trigger, not from the steps returning quickly. Measure the trigger against an immediate run to know.
  • What tells you, without reading any documentation, which model a given entry point uses?
    Point it at a large source and time a step that must touch every row. Under an immediate model that step's own line takes proportional time. Under plan-then-run it returns in microseconds and the time appears later at the trigger. A second tell: ask for the row count or print the object and watch whether work starts.

saying these in an interview costs you the question

  • The last call must be doing the expensive work itself
  • The first three steps were fast because the data was cached
  • Each of the three steps already produced a table in memory
  • Per-step timings from a deferred chain are a usable profile
  • The four minutes prove deferring is slower than running each step
open as a page

An exact median over a 40 GB file cannot be answered from summaries of pieces, but a total can — why?

level: juniorimportance: must knowfreq 62%

basics

~20 s

A total is determined by the records already read, but an exact median is not: a record still unread can move the middle value. Only the total has a combine step that absorbs whatever arrives later.

open as a page

A 2 GB text file of records is loaded whole and the process now holds far more than 2 GB. What sets that multiple?

level: juniorimportance: must knowfreq 70%

basics

~20 s

The in-memory representation of the columns sets the multiple, not the row count. Values held as one runtime object per row, header and reference included, cost several times their text; fixed-width typed columns can cost less than the file did.

open as a page

A large file is read in pieces of unequal row counts and each piece's mean is recorded. Why is the average of those piece means not the file's mean?

level: juniorimportance: must knowfreq 70%

basics

~20 s

Averaging piece means gives every piece one vote while the pieces hold different numbers of rows, so the answer drifts toward whatever the small pieces contained. Carry a running total and a running count instead, and divide once at the end.

open as a page

What does a recorded plan buy on one machine, where there is no network round trip to avoid?

level: middleimportance: must knowfreq 58%

basics

~20 s

Three local savings: columns the plan never references are never read or decoded, a row condition can be applied during the read instead of after it, and consecutive element-wise steps can collapse into one traversal with no intermediate buffer.

open as a page

A single-process data job 'ran out of memory' overnight; which distinct failures does that phrase cover, and how do you tell them apart?

level: middleimportance: must knowfreq 58%

basics

~20 s

The phrase covers at least four failures: a refused allocation inside the program, the process killed from outside for total resident size, a local disk filling with intermediates, and paging that ruined the wall clock. Different evidence, different step.

open as a page

A step whose finished result occupies 6 GB is killed on a 16 GB machine. What actually had to fit?

level: middleimportance: must knowfreq 58%

basics

~20 s

The peak had to fit, not the result. During the step the input is still referenced while the output and any intermediate are being allocated, so the worst instant can hold several copies of the data at once.

open as a page

A recorded plan is triggered twice for two different outputs and the second run takes as long as the first — why?

level: middleimportance: should knowfreq 46%

basics

~20 s

Nothing was retained. A recorded plan describes work rather than holding a result, so a second trigger replays the whole chain from the source unless the shared part was explicitly kept — held in memory, or written out and read back.

open as a page

An exact distinct count is needed over a file far larger than memory. What are the three escapes, and what does each charge?

level: middleimportance: should knowfreq 55%

basics

~20 s

Three: read the input more than once while holding intermediates on local disk, accept an answer with a stated error bound, or spend one cheap extra traversal so the second fits. They charge volume, exactness, and an extra read.

open as a page

A large input is folded a piece at a time and each piece reports its own variance. Which quantities must a piece carry instead so the pieces combine exactly?

level: middleimportance: should knowfreq 45%

basics

~20 s

Each piece must carry three things: its count, its mean, and its summed squared deviations from that mean. Variances cannot be averaged because each is measured around a different origin, and the merge has to add back the gap between those origins.

open as a page

A deferred pipeline fails at the trigger with an error about a step written twenty lines earlier — why there, and what would have caught it sooner?

level: seniorimportance: should knowfreq 52%

basics

~20 s

The step never ran where it was written — it ran at the trigger, so the report points there. How avoidable that is depends on the design: some planners resolve the schema while building and reject an unknown column immediately.

open as a page

A distinct-customer figure must tie out exactly for an audit, and the job must finish inside a 20-minute window. Which escape does each requirement rule out?

level: seniorimportance: should knowfreq 44%

basics

~20 s

Exactness rules out the bounded-error answer; the wall-clock budget rules out both escapes that re-read the input. Together they leave nothing, so a requirement, the hardware, or the amount of data the operation must see has to change.

open as a page

A single-machine data job is called slow with no further detail; which measurements name the resource that ran out and the step?

level: seniorimportance: should knowfreq 50%

basics

~20 s

Three numbers per step, not per run: elapsed time, processor time summed across cores, and peak resident bytes. The gap between elapsed and processor time separates waiting from computing; one busy core beside idle ones separates a design's default from a bottleneck.

open as a page

A large table is released and live data is now tiny, yet the process's resident size has not dropped. Is that a leak?

level: seniorimportance: should knowfreq 38%

basics

~20 s

Usually not. Resident size is a high-water mark: a general-purpose allocator keeps freed pages for reuse rather than handing them back, so the process stays near its peak with almost nothing live. A leak shows as live bytes that keep climbing.

open as a page

A per-column memory total reports 400 MB for a table while the process holds several gigabytes. What did that total leave out?

level: seniorimportance: should knowfreq 45%

basics

~20 s

A shallow total counts each column's own array of fixed-width slots. Where a slot is a reference, the object it points at - header and characters - is never counted, so text-heavy tables are understated by most of their real cost.

open as a page

A pass reads a large file a piece at a time yet still exhausts memory totalling amounts per customer, and halving the piece size does not help. What is the peak footprint proportional to?

level: seniorimportance: should knowfreq 58%

basics

~20 s

Live bytes at the peak are roughly one piece plus the fold's accumulator. Halving the piece shrinks only the first term. Here the accumulator holds one entry per distinct customer, so it is proportional to the key's cardinality, which the piece size never touches.

open as a page

Before a team rewrites a working single-machine data job, which two costs are routinely left out of the comparison?

level: principalimportance: should knowfreq 42%

basics

~20 s

Two: the hourly price of renting a machine several sizes larger, weighed against the engineering days the rewrite genuinely takes; and the standing cost of two implementations of the same logic, which must be changed twice and reconciled for as long as both exist.

open as a page

An exact median over a column that will not fit is computed in two traversals of the input. What does the first traversal buy?

level: seniorimportance: nice to knowfreq 38%

basics

~20 s

The first traversal computes a small fixed-size summary — counts per value range — identifying the one narrow range that holds the answer. The second traversal retains only that range, which fits, so an exact answer comes from a bounded footprint.

open as a page