Two nightly datasets are joined on customer id; what must be true of both writes for that join to move no records?
answer
- three things must match, not one
- same column on both sides
- same function from key to piece number
- the count lives inside the mapping
- the writer pays the move once
basics
~20 sBoth 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 sThree 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
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.
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.
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.
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