When a large input is matched against a tiny one, why does copying the tiny input to every worker avoid a redistribution?
answer
- two inputs, only one moves
- the larger side stays put
- one full copy per worker
- local keyed lookup, streamed match
basics
~20 sMatching records on a key normally forces every worker to send each record to whichever worker owns that key. Copying the whole small input to every machine instead lets the large input be matched where it already sits.
solid answer
~40 sA match on a key normally needs a redistribution: every worker sends each record it holds to whichever worker will handle that record's key, so equal keys meet. That puts both inputs on the network. If one input is small enough, you can instead place a complete copy of it on every worker - a worker being one process on one machine that holds a slice of the job's memory and runs some of its pieces - and each worker then matches its own share of the larger input against the local copy. The larger input never leaves the machine that read it. Who makes that choice varies: a declarative surface may pick it from size estimates, while a program written as per-record functions needs the author to register the copy explicitly.
go deeper
Recall the shape: one input is small enough to copy whole to every machine, so the big one never moves. Be able to say which side is copied and which stays where it was read.
Explain the alternative being avoided - every record sent to the worker that owns its key - and why matching against a resident local copy lets a worker emit results without waiting for other producers.
Show that the copy is paid once per worker, and that whether the runtime selects this shape at all depends on the surface you wrote. Name what you would measure before relying on it.
Frame it as trading replicated memory for network volume, and say at what cluster width and input size that trade stops paying - including what it does to the other jobs sharing the same machines.
## The movement this avoids A step that matches two inputs on a key can only run where the matching records sit together. In a cluster they do not start that way. Each **worker** - one process on one machine that holds a slice of the job's memory and runs some of its pieces - reads a **piece** of each input, one contiguous share it processes on its own, and nothing guarantees that the records which must meet are on the same machine. The general remedy is a **redistribution**: every worker sends each record it holds to whichever worker will handle that record's key, so equal keys always land together. It is correct for any pair of inputs, and it is expensive, because it puts *both* inputs on the network - and, where the job is finite and cut at a waiting line, on the producers' local disks as well. ## What copying the smaller input does instead If one of the two inputs is small enough to sit in a single worker's memory with room to spare, there is a second shape: send **the whole of the smaller input to every worker** and leave the larger one where it is. Each worker loads that complete copy into a keyed lookup in its own memory - a structure that answers "give me the records carrying this key" - and streams its own piece of the larger input past it, emitting matches as they are found and discarding each record once it has been matched. Three consequences follow: - **The larger input does not move at all.** It is read where it already lives and matched in place, so the dominant byte volume never touches the network. - **This match needs no waiting line.** A redistribution in a finite job typically introduces a line across the job where no downstream worker may compute until every upstream producer has finished; a worker matching against a local copy already holds everything it needs, so it can produce output immediately. Other steps in the same job may still impose such a line. - **The smaller input is now replicated.** One copy exists per worker, at the same moment, and that replication is the whole of what the technique costs. ## The two shapes side by side | | Redistribute both inputs | Copy the smaller input everywhere | |---|---|---| | Bytes on the network | Both inputs, once each | The smaller input, once per worker | | The larger input | Moves | Stays where it was read | | Memory held per worker | One incoming share | A full copy of the smaller input | | Correct when | Always | Only while the copy fits, with headroom | | Sensitive to | Uneven key distribution | Cluster width, and the copy's in-memory size | ## Who decides that it happens This is the part a candidate who has used one engine will overstate, because the three common answers really do differ: 1. **On a declarative surface** - where you describe the result you want and the runtime picks a physical shape - the runtime may choose this itself, weighing an estimate of the smaller input's size against a limit it is configured with. The estimate may come from statistics recorded alongside the data, from the metadata of the files it will read, or from a measurement taken after an earlier step finished. 2. **In a program written as per-record functions**, nothing estimates anything and nothing is rewritten. The author decides and registers the value to be replicated; the runtime executes what is written. 3. **In the oldest two-phase, disk-to-disk model** of this class, the copy is arranged before the job runs by putting the small input where every unit of work can read it locally, rather than being selected by a planner at all. So "the engine notices and does this for you" is true for some engines on some surfaces, and plainly false elsewhere. The mechanism is identical in all three; only the selection differs. ## What it does not fix - It does nothing when **both** inputs are large - that is a different strategy and a different subject. - It does not remove the requirement that both inputs carry the matching key. - It does not help when one key holds an enormous share of the larger input: the copy is fine, the *work* is still uneven, and that unevenness (skew) is its own subject. - It is not the same thing as having written both inputs in advance cut the same way on the same key, which avoids movement without copying anything. ## The sentence to say out loud "Matching on a key normally drives both inputs across the network so equal keys meet. If one side is small, you copy that side in full to every worker and match the large side in place - trading one resident copy per worker against moving the large input at all."
- Does copying the small input remove the waiting line where downstream compute halts for upstream producers?For that match, yes: each worker already holds everything it needs, so it emits results as it reads its share of the larger input. It does not remove every such line from the job - an aggregate or a sort later in the same job can still impose one. The copy removes this exchange, not all synchronisation.
- Which input is held in memory and which is streamed, and why does the order matter?The copied smaller input is held in a keyed lookup; the larger input is streamed past it record by record and released as it goes. Reversed, the worker would have to hold the larger input resident, which is precisely the cost the technique exists to avoid.
- Both inputs are moderate, not tiny. Is copying still the obvious choice?No. The copy is only free of the large input's movement while it comfortably fits on every worker at once. Once the smaller input is merely smaller rather than small, replicating it across a wide cluster can move more bytes than redistributing both would have.
saying these in an interview costs you the question
- Thinks the larger input is also copied to every worker.
- Believes copying the smaller input is always cheaper than redistributing both.
- Assumes every engine picks this shape automatically from statistics.
- Says the copy costs the small input's size once for the whole job.
- Confuses it with two inputs already written cut the same way on the key.