skip to content

Write and Fetch Sides

One side writes its records into a bucket per destination and the other collects its share from every producer: what each side does, and what it costs when a producer's output is gone.

on this pageshow

questions

4

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%

answer

  1. two halves, one exchange
  2. route first, then collect
  3. one bucket per destination, not per key
  4. every collector touches every producer

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.

solid answer

~50 s

A redistribution — every worker sending each record it holds to whichever worker will handle that record's key — has two halves. On the **producing side**, a unit walks its piece of the input, applies a routing rule to each record's key to get a destination number, and accumulates records into one bucket per destination, grouped so each destination's share sits together. On the **collecting side**, each unit owns one destination number and gathers that numbered share from *every* producing unit that ran, merging the arrivals into its own input. The halves are not symmetric: a producer makes as many buckets as there are collectors, a collector gathers as many shares as there were producers. Where the shares are written down before anything is collected, gathering can happen after the producer has exited; where the runtime instead pushes each record downstream as it is made, the same routing happens but the bucket is never stored bytes.

go deeper

for a junior

Recall the two halves: one side routes each record into a bucket for its destination, the other side gathers one numbered bucket from every producer that ran.

for a middle

Walk through both halves concretely — the routing rule applied per record, buckets grouped by destination number, and a collector merging one numbered share from every producing unit before it can finish a group.

for a senior

Say which regime you are describing. Shares written down and requested afterwards, or records pushed to a consumer that is already running, differ in when both sides must be alive and in whether a share can ever be read twice.

for a principal

The angle is what an exchange's shape commits the platform to: a pushed hand-off holds both steps' capacity at once, while a written-down one frees producing capacity at the price of intermediate bytes that live on machines you may lose.

## The moment that forces a movement Most steps in a distributed job are local. A filter, a projection or a per-record transformation runs on whatever records a worker already holds, and no bytes cross the network. Some steps cannot work that way. Grouping by a key, matching two inputs on a key, or re-cutting the input into a different number of pieces all need records that other workers are holding. That forces a **redistribution**: every worker sends each record it holds to whichever worker will handle that record's key, so all of a key's records end up in one place. This is the exchange that the batch lineage of these engines calls a shuffle, and it is the dominant cost of most distributed jobs. It is built as two halves, and an interviewer expects you to describe both: a **producing side** that decides where each record goes and packages it, and a **collecting side** that assembles its own share. ## The producing side: one bucket per destination A producing unit — one worker computing one piece of the input once — does four things with the records it makes: 1. **Decides a destination for each record.** It applies a *routing rule* to the record's key: a rule that names exactly one destination for that key, typically a hash of the key reduced to the destination count, or a set of boundaries over the key range. The rule must be the same function over the same destination count on every producing unit. If two producers disagree, a key's records land in two different places and every group over that key is quietly partial. 2. **Accumulates records into buckets.** There is one bucket per *destination* — not one per key, and not one per producer. Many keys share a bucket; one key never spans two. 3. **Arranges the buckets so one share can be handed over alone.** Records are grouped or sorted by destination number so destination N's records sit contiguously, with enough bookkeeping to find where share N begins. Designs differ here more than anywhere else in the exchange: one output plus an index of per-destination offsets, one output per destination, or nothing on storage at all when records are pushed straight out. 4. **Makes the share available** — by writing it to the producing machine's local disks, where it waits until a collector asks for it, or by sending it immediately to a consumer that is already running. The width of the *next* step is therefore what sets how many buckets each producing unit creates. ## The collecting side: one bucket number, every producer Each collecting unit owns one destination number, and its input is the union of that numbered share from **every producing unit that ran**. It obtains those shares, merges them, and then does the step's real work — folding a group, building a keyed lookup, matching runs of equal keys. Two properties of this half are what interviewers listen for: - **It gathers from all producers, not from one.** Any producing unit could have held records for any key, so a collector holding 199 of 200 shares does not yet have a complete group for a final result: the missing share may carry records for every key it owns. In a continuous job that emits an updated running result rather than one final answer, the collector still emits — but what it emits is explicitly provisional. - **Its work scales with the producing width**, while a producer's work scales with the collecting width. The transfer count is the product of the two, which is why widening both sides at once is so much more expensive than widening either. ## The same routing, two very different regimes The engines in this class genuinely disagree about what sits between the halves, so say which regime you mean: | | shares written down first | records pushed as produced | |---|---|---| | what a bucket is | bytes on the producing machine's local disk | a small send buffer per destination | | who starts the transfer | the collecting unit asks for its numbered share | the producing unit sends downstream as it goes | | liveness | the collector can run after the producer has exited | both sides are alive for the life of the exchange | | reading a share twice | possible while those bytes survive | impossible — nothing was kept | A third common arrangement runs a continuous computation as a rapid succession of small finite jobs; each little job performs the written-down version at its own small scale. ## Two traps in the vocabulary - The producer's bucketed output landing on local disk is the **normal mechanics of the exchange**, not a symptom of pressure. A working set written out because it does not fit a worker's memory budget is spilling — a different subject with different remedies, even though both show up as bytes on the same disks. - The **sort** here is *into buckets by destination*. It is not the sort that produces an output ordered end to end across the whole job, which needs boundaries drawn over the key range rather than a hash.

  • Why must every producing unit use the same routing rule and the same destination count?
    Because the whole point is that equal keys meet. If one producer maps a key to destination 7 and another maps it to destination 12, that key's records arrive at two collectors, each of which sees a partial group and emits a wrong answer without any error. The rule and the count are fixed before the producing step runs, for all of it.
  • Does a collecting unit know in advance how many shares it should receive?
    It knows how many producing units the step ran, so it knows how many shares to expect, and that count is what tells it when its input is complete. Some designs skip the transfer entirely for a share a producer has nothing in, so the number of transfers actually made can be lower than the number of producers.
  • What decides how many buckets a producing unit creates?
    The width of the collecting step — the number of destinations — and nothing else. It is not the number of distinct keys it saw, nor the number of producers. Changing that width after the producing step has run means the routing has to be redone, which is another movement.

A regional sorting office does not decide which letters it will receive. It takes the day's post it happens to hold, reads each address, and drops the letter into one sack per destination depot. Every depot then collects its own sack from every sorting office in the country, and only once the last sack is in can it claim it has the whole day's post for its area.

saying these in an interview costs you the question

  • Says a grouping can route records round-robin, so equal keys never meet
  • Thinks a collector pulls from one producer rather than from all of them
  • Believes a producer writes one bucket per distinct key it saw
  • Calls the producer's bucketed output spilling caused by memory pressure
  • Assumes every engine writes the shares down before anything is collected
  • Thinks the collecting side chooses which records it will be sent
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

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