skip to content

Arranging to Not Move

Two datasets already divided the same way on the same key can be combined where they sit; what has to hold for the engine to believe it, and what it costs to keep that true.

on this pageshow

questions

4

Two nightly datasets are joined on customer id; what must be true of both writes for that join to move no records?

level: middleimportance: should knowfreq 45%

answer

  1. three things must match, not one
  2. same column on both sides
  3. same function from key to piece number
  4. the count lives inside the mapping
  5. the writer pays the move once

basics

~20 s

Both writes must have used the same key, the same function from key to piece number, and piece counts the engine can relate - normally the same count. Only then does every matching pair already sit on one worker.

solid answer

~50 s

Three things have to line up, not one. Both datasets must be divided on the **same column** - the one the join matches on; both must have used the **same division rule**, the function from a record's key to the number of the piece it belongs in; and both must have been cut into piece counts the engine can relate, which normally means the **same count**. That third condition surprises people: almost every such rule folds the count into the mapping, usually a hash reduced modulo the number of pieces, so changing the count moves every key. When all three hold, all the records for a key already sit on one worker thread on both sides, and the join is a step each worker finishes alone rather than a wide step - one needing records currently held by every other worker.

go deeper

for a junior

Recall that a join only avoids the network when the matching records for a key are already on the same worker, and that this depends on how both datasets were written, not on how the join is written.

for a middle

Explain all three conditions and why the count is one of them: the rule reduces a hash by the number of pieces, so changing the count relocates essentially every key. Say that a shared column alone buys nothing.

for a senior

Show that the saving is moved rather than removed - the writer pays a redistribution and freezes a count into the data - and judge when repeated joins on that key earn it back.

for a principal

Frame the count as a cross-team contract on a shared dataset. Weigh one arranged key against the variety of keys real consumers join on, and against who is obliged to maintain it.

## What 'not moving' means here A distributed run divides its input into **pieces**: one slice of the stored data that a single worker thread reads and processes from start to finish. A join matches records by key, so a worker can finish one alone only if, for every key it holds on one side, it also holds every record carrying that key on the other side. When that is not true, the join is a **wide step** - a step a worker cannot finish from the records already in its hands, because it needs records sitting on every other worker - and records are redistributed across the network by key until matching ones meet. **Co-divided inputs** are two datasets written so that the matching already holds: the same key, the same division rule, and piece counts the engine can relate. Nothing moves at read time because the arrangement was paid for earlier, at write time. Note that this is a different thing from grouping stored files into directories named for one column's value; that arrangement lets a filter drop whole directories, but it says nothing about where any individual key lands relative to another dataset. ## The three conditions | Condition | What it means | What happens if it differs | |---|---|---| | Same key | Both sides divided on the column, or the ordered list of columns, the join matches on | the two layouts have no relationship at all; every key scatters | | Same division rule | The same function from a record's key to the number of its piece | one key lands in different piece numbers on the two sides | | Relatable counts | Normally the same number of pieces on both sides | the mapping shifts for essentially every key | A few details sit inside those rows and are worth saying out loud: - With a composite key, the **order and the exact set of columns** are part of the rule. Dividing one side on (customer id, region) and the other on (region, customer id) is not the same rule, and dividing on a superset is not either. - The rule has to agree on **type handling**: the same logical key written as a string on one side and an integer on the other hashes differently. - Records with a missing key are placed by whatever the rule does with an absent value, and that is a place two stacks quietly disagree. ## Why the count is inside the rule The usual rule is a hash of the key reduced modulo the number of pieces. That last step is what makes the count a condition rather than a detail: if one side has 200 pieces and the other 400, a key at piece 17 on the small side may be at 17 or 217 on the large side. Some engines exploit exactly that relationship and pair up pieces when one count is an exact multiple of the other; others accept only equality and redistribute both sides. That is why the honest phrasing is 'counts the engine can relate', and why in practice teams write both sides at one agreed count. ## Who pays, and when The saving is not free; it is moved. The producing job pays it: 1. Writing a dataset in this arrangement is itself a wide step - the writer must send each record to the piece its key belongs in. 2. The piece count is frozen into the stored data and becomes a contract that later writers and readers both depend on. 3. If the key is uneven, some output pieces end up far larger than the rest. One key holding more records than a single worker can take is a separate subject with its own remedies, and this arrangement does nothing to help it. The arrangement earns its money when the same two datasets are joined on the same key repeatedly - a nightly job, or many jobs, over data written once. ## Finite jobs and continuous jobs are not the same case In a **finite job** - a run over an input that ends - the pieces are usually derived from the stored bytes, so a layout recorded in storage can be adopted directly by the run that reads it. In a **continuous job** - a run over an input with no end - nothing about the input can be measured in advance, so the division is a **declared operator width** the author states, and it stands until the job is restarted. Two inputs are arranged to meet there when they enter at the same width, with the same key extracted and the same rule applied; a restart at a different width redistributes everything, and no property of the stored files prevents that. ## Being right is not yet enough One last thing separates the theory from the saving. Two datasets can satisfy all three conditions by accident and still be redistributed, because the run has to *know* the layout holds, not merely benefit from it. What the engine can prove, and how it comes to know, is the harder half of this subject.

  • One side has twice as many pieces as the other. Is anything salvageable?
    Sometimes. When one count is an exact multiple of the other, each small-side piece corresponds to a fixed set of large-side pieces, and some engines pair them up and still avoid the move. Others treat unequal counts as no match and redistribute both sides. The portable fix is to agree one count and rewrite the side that disagrees.
  • Both sides are divided on the join key, but by different hash functions. Does that help?
    No. The key is right and the mapping is not, so a given key sits in unrelated piece numbers on the two sides and every record still has to travel. You have paid the write-time cost of an arranged layout and bought nothing, which is why the rule - not just the column - has to be part of the agreement.
  • Does sorting each side by the join key remove the need to move records?
    Not by itself. Sorting rearranges records inside each piece; it does not change which piece a key is in, so matching records can still be on different workers. Sorting can make the join cheaper once the records are already together, but the together part is what the division buys.

Two warehouses that independently agreed to file every customer folder on the shelf numbered by the same rule. A clerk standing at shelf 17 has both halves of every customer whose folder belongs there, and never walks. Renumber one warehouse's shelves and the agreement is gone, even though both are still filed by customer.

saying these in an interview costs you the question

  • Says joining on the same column is enough, whatever the layout
  • Assumes any two hash-divided datasets line up automatically
  • Thinks the piece count is a tuning detail, not a matching condition
  • Believes sorting each side by the key removes the move
  • Treats the arrangement as free, forgetting the writer pays it
  • Confuses it with grouping files into directories named by a column
open as a page

A join between two tables stopped avoiding redistribution six months after both were written to match - what eroded the match?

level: seniorimportance: should knowfreq 33%

basics

~20 s

Later writes that did not follow the agreement: a backfill or a second pipeline appending files divided some other way, a piece count re-tuned on one side, or a key column quietly redefined. Nothing fails; the saving just disappears.

open as a page

Both join inputs were written with an identical division on the join key, yet the run still redistributes both sides - why?

level: seniorimportance: should knowfreq 42%

basics

~20 s

Because the run must know the layout holds, not merely benefit from it. Unless the key, rule and piece count are recorded where the plan reads them and survive every earlier step, it moves the records.

open as a page

Forty jobs join one shared table on the same key - should the platform require every write to preserve a matching division, and what does that cost?

level: principalimportance: nice to knowfreq 26%

basics

~20 s

Only if that key dominates the joins and the table is read far more often than written. The standard freezes a piece count into shared data, taxes every writer, and needs an owner and a check to survive.

open as a page