A job's 30 GB input is stored in a form that yields exactly one piece — what does adding machines buy?
answer
- capacity is not divisibility
- one unit, one lane, whole run
- idle from the first second
- divide the form, or redistribute after
basics
~20 sNo extra concurrency. One piece is read by one worker thread from start to finish, so the run's ceiling is one lane and the rest of the cluster idles no matter how many machines are added.
solid answer
~50 sIt buys nothing that this run can use. A **piece of the input** is read start to finish by a single **worker thread**, so a run whose stored form yields one piece has a parallelism ceiling of one: 30 GB go through one lane while every other lane in the cluster sits empty. Two things change that, and neither is capacity. First, change the stored form so it divides — then the derivation yields many pieces and many lanes fill. Second, accept the serial read and redistribute records afterwards, so every step past the redistribution runs at a larger count; the read stays a bottleneck but the rest of the job does not. The clue on a running job is one busy lane against hundreds of idle ones from the very first second — which is a different picture from a cluster that starts wide and narrows to a single slow lane at the end.
go deeper
Remember the shape rather than the remedies: one unit of work means one busy lane, and cluster size cannot change that. Saying so confidently is already a good answer at this level.
Walk through why each capacity remedy fails, then give the two that work: make the stored form divide, or redistribute after the read so later steps run wide. Note that the second leaves the read itself untouched.
Diagnose before prescribing. Distinguish one-unit-from-the-start from one-heavy-unit-at-the-end, say which evidence separates them, and weigh rewriting the stored data against paying for a redistribution on every run.
Treat it as a write-side standard, not a job-side fix: inputs produced in an undividable form tax every future reader. The call is whether to mandate a divisible form at the producer and who absorbs the rewrite.
## A ceiling of one A **piece of the input** is one slice of the stored data that a single **worker thread** reads and processes from beginning to end; a worker thread is one lane inside a **worker process**, which is one operating-system process on one machine. Because a piece is handed out once and is not subdivided mid-flight, the number of pieces is the ceiling on how many lanes can be busy. When the stored form yields exactly one piece, that ceiling is one — and a ceiling of one is indifferent to the size of the cluster below it. So the answer to "what do more machines buy" is: more lanes that stay empty. Thirty gigabytes stream through a single lane; the remaining lanes have nothing to be given. The run costs more and finishes at the same point. ## Why capacity cannot rescue it It is worth being precise about each of the plausible-sounding remedies, because interviewers probe them one at a time: - **More machines** raise how many worker processes can exist. Processes are containers for lanes, and lanes need pieces. - **More lanes per process** raise the worker-thread count, which was never the binding number here. - **Copies of the input on more machines** do not help either: every machine would read the same records and produce the same output. Concurrency requires *different* pieces, not duplicate reads. - **A larger or faster worker process** can change wall-clock time for the single lane doing the work — a bigger machine reads and decodes faster — but it does not change the ceiling, and how to size that process belongs to the subject on shaping a worker. The distinction worth stating out loud is between *capacity* and *divisibility*. Adding capacity to an undivided problem is like adding checkout counters to a shop with one queue that cannot be split: the counters are real, and no one can reach them. ## The two things that do change the picture 1. **Make the stored form divide.** Rewriting the input so that it can be decoded from many offsets — many moderate files rather than one undividable object — lets the derivation report many pieces, and the ceiling rises with them. This is the durable fix and it is paid for once, by whoever owns the write. *Why* a given stored form resists division is owned by the sibling leaf on the layout a reader inherits (`Layout the Reader Inherits`); this leaf owns only its consequence for the ceiling. 2. **Redistribute after the read.** A step that a worker cannot finish from the records it already holds, because it needs records currently sitting on every other worker, is a **wide step**, and it can emit into a larger number of downstream pieces. The read stays one lane wide, but everything after the redistribution runs wide. This is worth doing when the work *after* the read dominates — parsing, joining, aggregating — and worthless when the read itself is the cost. It is not free: every record crosses the network once, and the arithmetic of that cost is its own subject. ## What it looks like on a running job | Symptom | Likely cause | Where it belongs | |---|---|---| | One lane busy from the first second, the rest never used | the stored form yielded one piece | this subject | | Many lanes busy, then one lane still running at the end | one piece far heavier than the others | the subject on skew and stragglers | | Lanes busy in waves with gaps between them | work arriving in groups separated by a redistribution | the subject on moving data between workers | The first row is the one this question is about, and it is diagnosable in seconds: count the units of work the run reports against the lanes available. If the first number is one, no amount of the second number matters. ## Does this depend on the runtime? The ceiling does not; the shape of the premise does. On a runtime that re-derives the division from the stored bytes on every finite run, the moment the stored form becomes divisible the next run widens by itself. On one that settles the split count at submission, the new form is picked up at the next submission. In a **continuous job** — a run over an input with no end — the same situation appears as a source that offers only one independent unit to read from: declaring a larger operator width does not divide that source, so the read stays single-laned while the wider operators downstream wait on it. In all three the lesson is identical: the count of independent units of work, not the capacity, is what is binding.
- Would redistributing records straight after the read make the job faster here?Only if the work after the read dominates. The read itself stays one lane wide, so a redistribution helps parsing, joining and aggregating run wide, and helps nothing if decoding the 30 GB is the cost. It also pays for every record to cross the network once.
- How do you tell this apart from one unusually heavy unit of work among many?Look at the start, not the end. Here one lane is busy and the rest are empty from the first second, and the run reports a single unit of work. With one heavy unit among many, the cluster starts wide and thins to a single straggler — a different subject with different remedies.
saying these in an interview costs you the question
- Suggests a bigger cluster as the fix for an undivided input.
- Thinks copying the input to every machine creates parallelism.
- Says raising lanes per worker process will occupy the idle capacity.
- Assumes the runtime will cut the single unit up by itself.
- Confuses this with one heavy unit among many finishing late.