A job filters a billion rows down to four million and returns them all to the submitting program — what fails first?
answer
- the failure is in one process
- result size, not input size
- count times per-record width
- more workers change nothing here
- write from the workers, return a path
basics
~20 sThe coordinating process runs out of memory. Every returned record leaves the workers and lands in that one process's heap on one machine, so the cost follows the result size, not the cluster size — adding worker processes does not help.
solid answer
~50 sBringing results back to one place means moving every returned record out of the worker processes and into the single coordinating process — the one process that plans the pieces, hands them out and tracks what finished. The workers were never under pressure here: they each handled a slice. The four million surviving rows, however, are gathered into one heap on one machine, and four million rows at a few hundred bytes each is gigabytes. The filter ratio is a trap: removing 99.6% of a billion rows still leaves four million. The remedies are to write the rows from the workers to the destination store and return a path and a count, to return a folded aggregate, or to return a bounded prefix. Engines differ in how this presents: some enforce a cap and refuse quickly, others trickle the result back until memory is gone.
go deeper
Recall that returning records to the submitting program pulls them all into one process on one machine, and that the size of the result — not the size of the input — is what decides whether that works.
Explain the arithmetic out loud: returned record count times per-record in-memory width against one heap, and say why more worker processes or a finer cut of the input change neither term.
Demonstrate the production instinct: name the remedy before the sizing knob, know that some runtimes refuse fast with a cap while others trickle the result back until the process dies, and recognise the interactive display that is this operation in disguise.
The tradeoff angle is where a result belongs. A platform that lets any job gather an unbounded result into one process has a single-machine bottleneck in a distributed system; an enforced cap costs some convenience and buys a failure that is immediate, cheap and legible.
## The operation that moves records into one process Most operations in a distributed job leave the data where it is: each worker process reads its own pieces, computes over them, and writes its own output. A small family of operations breaks that rule. They move every returned record out of the worker processes and into the single **coordinating process** — the one process in the job that plans the pieces, hands them out, tracks what finished, and receives anything the program asks to bring back to one place — so the submitted program can hold the result as ordinary in-memory values. That family includes returning a full result set to the program, converting a distributed result into a single-machine table or array, and any interactive display that materialises the whole result rather than a bounded prefix. They share one cost shape and one failure. ## Why the workers survive and the one process does not The workers divided the billion input rows by their count: a job with two hundred worker processes gave each of them roughly five million input rows to scan, and the filter is cheap and streaming. The four million survivors are a different matter, because they do not stay divided. Their total in-memory size is: - **returned record count** multiplied by **per-record in-memory footprint**, held in **one heap** on **one machine**. The count is the number everyone quotes and the width is the number that decides it. Forty narrow integer columns and forty columns carrying free text are the same four million rows and two very different totals. Per-record footprint also varies several-fold between engines: some hold records in a compact binary layout the runtime manages itself, others hold ordinary language objects with per-field object overhead, and the same result set can differ by a factor of several between two runtimes of this class. The filter ratio is the trap the question is built on. Removing 99.6% of the input sounds decisive and is irrelevant: what matters is the absolute size of what survives, measured against one process's memory. ## How the failure presents, and what varies | runtime or deployment behaviour | what you observe | |---|---| | the result is fetched in chunks as pieces finish | memory climbs steadily, collection pauses lengthen, the process slows badly and then dies | | the runtime enforces a cap on returned size | a fast, explicit refusal naming a limit — the kinder failure, and the one you should want | | the coordinating process runs outside the cluster on the submitting machine | the ceiling is that machine's memory, often far lower than a cluster machine's | | a managed compute service, where you are shown no machine | the limit still exists and is the provider's, so you meet it as a service error rather than a heap dump | On any of them the diagnosis is the same, and so is the thing that does **not** help: more worker processes, more worker memory, or a finer cut of the input. None of those change the size of what is gathered into one place. ## The remedies, in the order to reach for them 1. **Do not return the rows.** Have the workers write them to the destination store in parallel and return a path and a count. This is the right answer for almost every production job; a result set large enough to worry about is a result set that belongs in storage. 2. **Return an aggregate instead of rows.** Fold in the workers so that one small row per group crosses back. A count, a sum per category or a set of percentiles is a few kilobytes regardless of how many rows fed it. 3. **Return a bounded prefix or a sample, with the bound in the code.** A limit that exists only as a habit is a limit that will be forgotten by the next person; make the runtime enforce it. 4. **If the rows genuinely must reach one place** — a report generator, a client library that only speaks single-machine tables — then page the result, write it to local disk as it arrives rather than accumulating it in the heap, and size that process for the result rather than for the job. ## What this is not It is not a worker memory problem; dividing one worker process's memory between the operators running in it is a separate subject with its own symptoms. It is not a network throughput problem in the first instance — the bytes usually arrive fine and then have nowhere to live. And it is not fixed by cutting the input into more pieces: the number of pieces changes how the work is divided, not how much comes back.
- Why does returning a grouped aggregate over the same four million rows normally succeed?Because the fold happens in the worker processes. Each worker reduces its own rows to one partial value per group, those partials are combined, and only one small row per group ever crosses back. The returned size is then bounded by the number of distinct groups rather than by the number of rows, which is usually smaller by orders of magnitude.
- How would you size the coordinating process for a job like this?From the bookkeeping for the number of pieces plus the total in-memory size of whatever is returned, never from the input volume. If that arithmetic produces a large number, treat it as a design signal rather than a sizing problem: a result that needs a specially enlarged process to be gathered in one place should almost always be written from the workers instead.
- The program needs the rows in one place only in order to write them as one file. Is that different?Yes, and it does not need this operation. Bringing the rows into the coordinating process and writing from there puts every byte through one heap; having a single worker process write the single file keeps the bytes in the cluster. How many output pieces a job should produce is its own subject, but the point here is that a one-file requirement is not a reason to return rows to the program.
A hundred stocktakers can count a warehouse in an afternoon because each takes an aisle. Asking each to report a total works fine. Asking each to drive every box they counted to the head office car park does not, and buying more stocktakers makes the car park no larger.
saying these in an interview costs you the question
- Adds worker processes to fix an out-of-memory in the coordinating process
- Assumes a filter removing most rows makes any result safe to return
- Reasons from the row count and ignores the per-row width
- Believes returned rows stream to the caller without being held anywhere
- Assumes every on-screen preview fetches only a bounded prefix
- Treats it as a network bandwidth problem rather than a memory one