skip to content

Moving Data Between Workers

A step that needs records held by other workers forces a redistribution - a shuffle. How that exchange is organised, what it truly costs, and how an author makes less of it happen.

on this pageshow

explore

questions

27

Two 500 GB inputs are matched on a shared account key, and neither fits on one worker - what happens to both inputs first?

level: juniorimportance: must knowfreq 72%

answer

  1. neither input can be copied everywhere
  2. matching is local or not at all
  3. one rule applied to both sides
  4. same destination count on both
  5. keys with no partner travel too

basics

~20 s

Both 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 s

Neither 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

for a junior

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.

for a middle

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.

for a senior

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.

for a principal

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.
open as a page

Why is the cost of redistributing data across a cluster counted in bytes moved rather than in records?

level: juniorimportance: must knowfreq 70%

basics

~20 s

A redistribution costs bytes: what crosses the network, what the producing side writes to local disk where output is materialised, and processor time spent encoding it. Record counts predict none of those unless every record is the same width.

open as a page

A job sorts each of its 200 pieces locally and writes them out; why is the combined result not ordered end to end?

level: juniorimportance: must knowfreq 70%

basics

~20 s

Sorting inside a piece orders only that piece. Unless the routing that filled the pieces gave each one a contiguous span of the key range, piece 3 can hold keys that belong before piece 2's, so reading the pieces in order interleaves ranges.

open as a page

Why does folding records together on the machine that produced them cut what a grouped aggregation sends across the network?

level: juniorimportance: must knowfreq 70%

basics

~20 s

Records sharing a grouping key are combined where they were read, so each producing machine sends one partial result per key instead of every record. The collecting side combines partials into the same final answer, over far less traffic.

open as a page

When a large input is matched against a tiny one, why does copying the tiny input to every worker avoid a redistribution?

level: juniorimportance: must knowfreq 72%

basics

~20 s

Matching 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.

open as a page

In a finite job, why must a step that needs records from other workers wait for every producing piece to finish?

level: juniorimportance: must knowfreq 72%

basics

~20 s

A step that needs records held by other workers cannot know it has them all until every piece of the producing step is done; computing sooner would publish a partial answer. That waiting line is a stage barrier.

open as a page

Two large inputs have been redistributed on the matching key - what two ways can a worker match the landed pieces?

level: middleimportance: must knowfreq 64%

basics

~10 s

Either load one landed piece into a keyed in-memory lookup and stream the other past it, or have both landed pieces arrive ordered by the key and advance through the two ordered runs together.

open as a page

How can a cluster produce an end-to-end ordered result without collecting every record onto one worker to sort?

level: middleimportance: must knowfreq 65%

basics

~20 s

Draw a sample of the ordering key, pick boundary values from it so each piece receives one contiguous span of the range, redistribute records by comparing the key against those boundaries, sort each piece in parallel, and read the pieces in index order.

open as a page

Two programs compute the same per-key total, but only one lets the runtime fold before the movement — what differs in how they are written?

level: middleimportance: must knowfreq 58%

basics

~20 s

One states the per-key combining operation as part of the grouping step, so partial results can be formed on each producing machine. The other asks for the whole group first and computes afterwards, which obliges the job to deliver every record before anything can be combined.

open as a page

Why is copying a 2 GB input to every worker more expensive as the cluster grows, when the input never changes size?

level: middleimportance: must knowfreq 60%

basics

~20 s

The copy is paid once per worker, not once per job: two hundred workers means 2 GB sent two hundred times and 2 GB resident on two hundred machines at once. Cost tracks cluster width.

open as a page

What does the producing side write, and the collecting side gather, when a step needs records that other workers hold?

level: middleimportance: must knowfreq 72%

basics

~20 s

Each producing unit routes every record into one bucket per downstream destination; each collecting unit gathers its own bucket number from every producer that ran. Where buckets are written down first, collection can outlive the producer.

open as a page

A 400-piece producing step feeds a 200-piece collecting step — how many bucket transfers does that all-to-all exchange involve?

level: juniorimportance: should knowfreq 58%

basics

~20 s

Up to eighty thousand — the product of the two widths. Each of the 400 producers holds one bucket for each of the 200 destinations, and each of the 200 collectors gathers one share from each of the 400 producers.

open as a page

When does compressing a step's output before it crosses the network make a redistribution slower rather than faster?

level: middleimportance: should knowfreq 48%

basics

~20 s

Compressing trades processor time for bytes on the wire, so it loses whenever the job is short of cores rather than bandwidth, or when the data is already compact, so the ratio approaches one and the encoding work buys almost nothing.

open as a page

Two requests resemble end-to-end ordering — the 100 largest values overall, and each key's records in time order; does either need range boundaries?

level: middleimportance: should knowfreq 45%

basics

~20 s

Neither does. The top few are found by keeping the best 100 inside each piece and merging those candidates, and ordering inside a key only needs equal keys to land together, which any routing on that key already gives. Range boundaries serve a complete result read in order.

open as a page

A job groups a billion records by a near-unique identifier, so why does folding locally before the movement save almost nothing?

level: middleimportance: should knowfreq 52%

basics

~20 s

Folding replaces the records sharing a key with one partial result, so its saving is the records-per-key ratio. When keys are near-unique that ratio is about one: every record becomes its own partial, and the traffic is unchanged while the bookkeeping is added.

open as a page

While the producing step of a finite job is still running, what may the consuming side already do, and what must it not?

level: middleimportance: should knowfreq 48%

basics

~20 s

Collecting bytes from producers that have already finished may overlap producers that are still running; computing on them may not, because the consumer's input is not complete until the last producing piece is done. Transfer overlaps the wait, compute does not.

open as a page

Why does a continuous job that hands each record downstream as it is produced never reach a line where compute waits for all producers?

level: middleimportance: should knowfreq 52%

basics

~20 s

Because its producers never finish. A line across a job says compute waits until every upstream piece is done, and over an endless input that moment never arrives, so operators must emit from partial input instead of waiting for completeness.

open as a page

One matching key holds a hundred times the records of any other - how does that change the memory each local matching strategy needs?

level: seniorimportance: should knowfreq 44%

basics

~20 s

The keyed lookup's resident set is the whole loaded side, whatever the distribution of keys inside it, so it is unchanged. The ordered walk's is the run sharing the key being matched, so one enormous key is exactly what inflates it.

open as a page

A landed piece is twice the memory a worker can give the match - what does each of the two local matching strategies do then?

level: seniorimportance: should knowfreq 50%

basics

~20 s

The keyed in-memory lookup cannot hold it: runtimes respond by cutting both landed pieces further on the key and matching sub-pair by sub-pair, by switching to the ordered walk, or by failing. The walk instead pays extra ordered passes over local disk.

open as a page

A step attaches a 40 KB document to every record just before a grouping redistributes 200 million of them. Why can that movement dwarf the compute on both sides of it, and what do you change?

level: seniorimportance: should knowfreq 52%

basics

~20 s

Width multiplies every byte of the movement after it: 200 million records at 40 KB is about 8 TB on the wire, against 12.8 GB for the grouped fields alone. The repair is to widen after the movement, not before.

open as a page

Range boundaries drawn from a sample leave one piece holding most of the records; what went wrong and what does it cost?

level: seniorimportance: should knowfreq 50%

basics

~20 s

The sample misrepresented the key distribution, so the cut points do not split the data evenly. The output is still correctly ordered; what suffers is the run, which is now bounded by the one oversized unit of work and its memory pressure.

open as a page

A per-key total runs through a step your runtime will not fold locally, so every record crosses the network — how do you fold it by hand, and what does that cost?

level: seniorimportance: should knowfreq 44%

basics

~20 s

Add an explicit first step that walks each piece of the input, accumulates per key in a bounded table, and emits one partial per key; then redistribute and combine the partials. It costs producer memory, a lookup per record, and a two-step shape a later reader must understand.

open as a page

A candidate input is 400 MB compressed in shared storage - what must you know before copying it to every worker?

level: seniorimportance: should knowfreq 50%

basics

~20 s

What it becomes in memory, not what it occupies on disk. Decoded into records and indexed for lookup, a compressed file commonly expands several-fold, and that expanded figure is what every worker must hold at once.

open as a page

An input copied to every worker turns out far larger than estimated - where does the job fail, and why on many workers at once?

level: seniorimportance: should knowfreq 46%

basics

~20 s

Wherever the copy is assembled or held: on engines that gather it centrally, the gathering process can exhaust its memory first, and otherwise every worker fails at nearly the same point, because all of them are doing the identical thing.

open as a page

A finite job has three lines where compute waits for every producer; doubling the worker count changes its wall clock barely at all. Why?

level: seniorimportance: should knowfreq 50%

basics

~20 s

Each line lifts when its slowest producing piece finishes, not when the average one does. If every piece was already running at once, extra workers add capacity nobody was waiting for, so the job's wall clock stays roughly the sum of three maxima.

open as a page

A machine holding a finished producing step's written buckets is lost while later steps run — what does recovering that cost?

level: seniorimportance: should knowfreq 54%

basics

~20 s

The producing work again, not just a transfer. Those buckets are intermediate bytes on one machine's local disks, usually in a single copy, so nothing can be re-read: engines that kept a description of how the piece was made re-run it.

open as a page

Where a producing operator pushes each record downstream as it is made, what replaces the bucket a consumer would otherwise collect?

level: seniorimportance: should knowfreq 46%

basics

~20 s

Nothing durable. The routing rule still picks one destination per record, but the bucket is only a send buffer flushed over a channel to a consumer that is already running. There is no stored share to request, and none to request twice.

open as a page