Two 500 GB inputs are matched on a shared account key, and neither fits on one worker - what happens to both inputs first?
answer
- neither input can be copied everywhere
- matching is local or not at all
- one rule applied to both sides
- same destination count on both
- keys with no partner travel too
basics
~20 sBoth inputs are cut on the matching key by the same routing rule into the same number of destinations and sent across the cluster, so records sharing a key value land on one worker, which then matches them locally.
solid answer
~50 sNeither input can be copied to every worker, so both have to move. Each worker applies the same routing rule to each record's matching key - a deterministic rule that maps a key value to exactly one destination - and sends the record there. Applied to both inputs over the same number of destinations, it guarantees that every record carrying a given key value, from either side, converges on one worker, which can then pair them without talking to anyone else. This is a redistribution: a shuffle, in the sense that every worker sends each record it holds to whichever worker will handle that record's key. Both inputs cross the network in full, including keys with no partner on the other side, because nothing can tell a partner is missing until both sides have landed.
go deeper
Recall that when neither input can be copied to every worker, both are cut on the matching key and sent across, so records sharing a key value land on one worker.
Explain why the same routing rule over the same number of destinations is required on both sides, and why records whose key has no partner still travel the whole way.
Say what this costs in practice: both inputs cross the network in full before a single pair exists, and the landed pieces still have to be matched by one of two local strategies.
Treat the recurring double movement as the thing a platform decision has to answer for when two large datasets are matched on a schedule; how the inputs are arranged when they are written is a separate design subject from how the match is expressed.
## The problem the movement solves Matching two datasets on a key means every record on one side must be compared with the records on the other side that carry the same key value. On one machine that is trivial, because everything is reachable. On a cluster it is not: the two inputs were produced at different times by different jobs, and the slice of each that any one worker holds has no relationship to the slice another worker holds. **A worker can only match records it holds.** So before a single pair can be formed, the records that belong together have to be brought together. There are only two ways to arrange that, and when both inputs are large only one of them is available: - **Copy one input in full to every worker**, so the other never leaves the machine already holding it. That works only while the copied input is small enough to sit in every worker's memory at once - a separate subject, and not an option when both sides are enormous. - **Move both inputs**, cutting each on the matching key so that the records which must meet arrive in the same place. ## The mechanism: one rule, applied identically to both sides The movement is a **redistribution** - a shuffle, in the precise sense that every worker sends each record it holds to whichever worker will handle that record's key. What makes it produce a matchable result is the **routing function**: a rule applied to a record's matching key value that names exactly one destination for it. Three properties do all the work. 1. **It reads the key value and nothing else.** Two records with the same matching key are indistinguishable to it, whichever input they came from. 2. **It is deterministic.** The same key value always yields the same destination, on this worker and on every other. 3. **It is applied to both inputs over the same number of destinations.** The destination count is part of the rule, not a separate tuning choice: a rule over two hundred destinations and a rule over three hundred are different rules and will disagree about where a key belongs. Given those three, every record carrying a given key value converges on one worker, which can then pair them entirely locally with no further communication. Note what is **not** promised. Nothing says each worker receives a similar amount of data, and nothing arranges the keys in order across workers. Co-location of equal keys is the whole of the guarantee. ## What travelling actually looks like The execution models in this class genuinely disagree here, and a sentence that fits one of them is routinely false of another. | execution regime | how the records cross | |---|---| | a finite job cut into steps at a line where downstream compute waits | the producing side writes its records into one bucket per destination on local disk, and the consuming side collects its own bucket from every producer that ran | | a continuous job that hands each record over as it is produced | records are pushed to their destination worker as they come into existence, with nothing materialised in between | | the older two-phase disk-to-disk model | the producing phase writes to local disk and the consuming phase collects, with the records for a key delivered grouped and already in key order | The internals of the writing and the collecting sides are their own subject. What matters for a match between two large inputs is the part all three regimes share: **both inputs cross the network in full.** ## Including the records that have no partner Routing happens before anything knows whether a partner exists. A producing worker holds a slice of one input and has no view of the other at all, so a key present on the left and absent on the right is routed, sent and collected exactly like any other, and only the receiving worker discovers there is nothing to pair it with. This surprises people who imagine the movement carries the matching rows. It does not; it carries all the rows, arranged. ## What lands, and what is still to do - each worker ends up holding two pieces - its share of the left input and its share of the right - that agree on the key space - those pieces are not necessarily similar in size to any other worker's - the records inside them are not necessarily in any order - **nothing has been matched yet** Pairing the two landed pieces is a second, purely local decision with two shapes: loading one of them into a keyed lookup in memory and streaming the other past it, or having both arrive ordered by the key and advancing through the two ordered runs together. The movement only guarantees that the worker now has everything it needs to answer that question without asking anybody else. ## What an interviewer is listening for That you say **both** inputs move; that you name the routing rule rather than a setting; that you know the destination count is part of that rule; and that you describe the crossing in terms of producers, consumers and keys rather than in one runtime's vocabulary, as though every engine of this class worked the way the one you know does.
- Why must both inputs be cut into the same number of destinations, not merely by the same rule?Because the destination count is part of the rule: it maps a key value to one of N places. Cut one input into two hundred destinations and the other into three hundred and the same key value names a different worker on each side, so the two records never meet. Equal keys converge only when the rule and the count agree.
- Do records whose key appears on only one side still cross the network?Yes. The destination is decided from the key value alone, and the producing worker holds a slice of one input with no view of the other, so it cannot tell that a partner is missing. Only the worker that receives both sides can. The movement is therefore paid for every record, matched or not.
saying these in an interview costs you the question
- Says only the larger of the two inputs has to be redistributed.
- Thinks the records that must meet already sit together by luck.
- Applies a different routing rule to each input and still expects keys to meet.
- Believes keys with no partner are filtered out before the movement.
- Assumes the two sides may use different numbers of destinations.
- Calls the movement free because it happens inside the cluster.