skip to content

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%

answer

  1. multiply the widths, do not add them
  2. every collector needs every producer
  3. P times C shares
  4. more shares, smaller each, same overhead

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.

solid answer

~50 s

The fan is the **product**, not the sum: 400 x 200 = 80,000 producer-to-collector shares. Every collector needs something from every producer, because any producer could have held records for any of the keys routed to that collector. The consequence people miss is what happens to the size of each share: the same total records are now cut into 80,000 pieces, so widening both sides tenfold gives eight million shares that are a hundredth the size each. Every transfer carries a fixed cost that has nothing to do with its records — a request, locating the share inside the producer's output, and bookkeeping on both sides — so a very wide exchange can be slow while moving a modest amount of data. Where a producer has nothing at all for a destination, some designs skip that transfer, so the product is an upper bound rather than an exact count.

code

python · 14 lines
python
producers = 400
consumers = 200

shares = producers * consumers          # 80_000 producer-to-collector shares
buckets_per_producer = consumers        # 200 buckets held by each producer
collections_per_consumer = producers    # 400 shares gathered by each collector

# doubling both widths multiplies the fan by four
wider = (2 * producers) * (2 * consumers)

# and each share shrinks as the fan grows, for the same payload
moved_bytes = 80 * 1024**3
avg_share = moved_bytes / shares        # about 1 MB
avg_share_wider = moved_bytes / ((10 * producers) * (10 * consumers))  # about 10 KB

go deeper

for a junior

Recall that the number of shares in an all-to-all exchange is the two widths multiplied, and be able to say why each producer holds one bucket per destination.

for a middle

Explain the consequence rather than the formula: as the fan grows the payload is cut into more and smaller shares, while the fixed cost of each transfer stays where it was.

for a senior

Recognise the symptom in a real job — a step slow at a width its byte volume does not justify — and separate a transfer-count problem from a payload problem before changing anything.

for a principal

The judgement is about defaults across a platform: a width that suits the largest job makes every small job pay a fan it never needed, and a single global setting cannot be right for both.

## The arithmetic Call the producing step's width **P** and the collecting step's width **C**. In a redistribution — every worker sending each record it holds to whichever worker will handle that record's key — each producing unit prepares one bucket per destination, and each collecting unit assembles one numbered share from every producer. So: - each producer holds **C** buckets: 200; - each collector performs **P** collections: 400; - the exchange as a whole involves **P x C** shares: 400 x 200 = **80,000**. That is the headline number, and it is a ceiling rather than an exact count: where a producer happens to have no records at all for a destination, several designs skip that transfer entirely, so a sparse key space can make the effective count noticeably lower. ## Why a product and not a sum All-to-all is not a figure of speech. A collector owns a set of keys, and records for those keys could have been produced anywhere — the input piece that happened to contain them is unrelated to the key they carry. So there is no producer a collector may skip on principle, and no collector a producer can ignore. Each pair is an independent share. That is why the two widths multiply: | change | shares in the exchange | buckets per producer | shares per collector | |---|---|---|---| | double the producing width | x2 | unchanged | x2 | | double the collecting width | x2 | x2 | unchanged | | double both | **x4** | x2 | x2 | A candidate who answers 600 has added the widths; a candidate who answers 400 or 200 has described one side's view of the exchange rather than the exchange. ## What a transfer costs beyond the records in it Each of the 80,000 shares carries overhead that does not shrink with its contents: - a request and its handling on both machines; - locating that share inside the producer's output — an offset lookup, or opening one more output; - per-share bookkeeping: which shares have arrived, which are outstanding, which failed and must be asked for again; - a bounded number of collections in flight at once on the collecting side, so a collector works through its producers in waves rather than all at once. Now put the two together. Suppose the step moves 80 GB. At 80,000 shares the average share is about 1 MB, and the fixed cost per transfer is noise. Widen both sides tenfold — 4,000 producers, 2,000 collectors — and the same 80 GB is cut into **8,000,000** shares averaging about 10 KB. Nothing about the data changed; the exchange now performs a hundred times as many transfers for the same payload, and the per-transfer overhead has stopped being noise. This is the shape behind a job that is slow at a width where the bytes alone do not explain it. ## How the regime changes what the product means The number is the same in every regime; what you are counting differs. - Where shares are **written down first**, the product counts transfers a collector requests, possibly long after the producer exited. - Where records are **pushed downstream as they are made**, the product counts logical channels held open simultaneously, one per pair, for the life of the exchange — plus a send buffer for each. Designs vary in whether each logical channel gets its own network connection or many are multiplexed onto one, which is exactly why the same width behaves differently on different engines. - Where a continuous computation runs as a **rapid succession of small finite jobs**, the fan is paid once per little job, so a high rate of little jobs multiplies the transfer count even when each one is tiny. ## What this arithmetic does not decide It tells you what a pair of widths costs in transfers. It does not tell you what the widths should be: how many pieces the input is cut into, and how wide the step after it runs, are chosen against parallelism, record counts and memory, and are a subject of their own. Nor is this the movement's cost in bytes on the wire — that is a separate calculation over record widths and encodings. The fan is the count; the bytes are the payload; a wide exchange can be expensive in either, and diagnosing one as the other is the common error.

  • The collecting step is widened from 200 to 800 while the producing step stays at 400. What changes?
    The fan goes from 80,000 to 320,000 shares. Each producer now holds 800 buckets instead of 200, so its output is cut four ways finer, and each collector still gathers 400 shares but each one is roughly a quarter the size. The producing side's per-bucket bookkeeping grows; the per-collector collection count does not.
  • Why can a wide exchange be slow even though the total bytes moved are unremarkable?
    Because the cost is partly per transfer, not per byte. Requests, share lookups and completion bookkeeping happen once per producer-collector pair, and the pair count grows as the product of the widths while the payload stays fixed. Past some width you are paying mostly for transfers that carry almost nothing.

saying these in an interview costs you the question

  • Says the transfer count is the sum of the two widths
  • Thinks doubling both widths doubles the number of shares
  • Assumes a tiny share costs nothing because its bytes are few
  • Confuses the bucket count with the number of distinct keys
  • Believes a collector can skip producers that are unlikely to hold its keys