A join between two tables stopped avoiding redistribution six months after both were written to match - what eroded the match?
answer
- every future write must honour it
- backfills and second producers
- a count re-tuned on one side
- nothing fails, it just moves again
- check the plan, not the clock
basics
~20 sLater 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.
solid answer
~40 sThe arrangement is a property of every write to both tables, forever, and only the first write was made under supervision. The usual erosions are a **backfill or repair job** that appends files without applying the same rule, a **second producer** that writes into the table a different way, a **count re-tuned** on one side as it grew, a **key column redefined** - widened, cast, made composite - so the recorded rule no longer describes it, and **maintenance that rewrites the table's files** without preserving placement, which is the table layer's concern rather than the engine's. The symptom is the quiet kind: results stay correct, the plan grows a redistribution back, and the job is simply slower and more expensive than its cost model assumed.
go deeper
Recall that an arrangement of stored data can be undone by later writes, and that when it is, the job still produces the right answers and simply does more work.
Name concrete causes - an appending backfill, a second producer, a re-tuned count, a redefined key - and explain why none of them raises an error anywhere.
Show how you would detect it deliberately: assert placement at write time, watch for the redistribution returning to the plan, and decide between rewriting the data and abandoning the arrangement.
Treat it as an ownership problem. Decide who is allowed to write the table, where the agreement is recorded, and whether the organisation can sustain the discipline the saving depends on.
## Why the match decays rather than breaks Co-divided inputs - two datasets written with the same key, the same division rule and the same piece count, so matching records already sit on one worker - are not a property you set once. They are a property that every future write to either dataset has to honour. The first load is written by whoever designed the arrangement and knows about it. The next hundred writes are made by people, jobs and maintenance processes that do not. The decay is silent by construction. A layout the run cannot rely on is not an error condition; it is simply a plan with a **wide step** in it - a step whose workers need records currently sitting on every other worker, so the records are redistributed first. The answers are identical. Only the duration and the bill change, and both change gradually enough to be attributed to data growth. ## The usual causes, roughly in order of how often they bite 1. **A backfill or repair job.** Someone reprocesses three months of history and appends the output. It is correct data, written the ordinary way, and it obeys no division rule. 2. **A second producer.** A different team starts writing into the same table from a different pipeline. Nobody told them the layout was load-bearing, because nothing in the table says so. 3. **A re-tuned count on one side.** The larger table grew, someone doubled its piece count to get more parallelism, and the two sides no longer relate. 4. **The key changed shape.** The key was widened, cast to another type, or became composite. The recorded rule describes the old key, not the one the join now matches on. 5. **Maintenance that rewrites files.** Consolidating or reorganising a table's files can move records between them. Whether placement survives that depends on the table layer, which owns this mechanism; the engine reading the result has no say. 6. **A restart at a different width.** Where the job is continuous, the division comes from a declared operator width rather than from stored bytes, so restarting either side at a new width redistributes regardless of what is on disk. ## What you actually observe | Symptom | What it tells you | |---|---| | Runtime creeping up over months with flat input size | something structural changed, not the data volume | | A redistribution now present in the plan before the join | the property is no longer established | | Network and temporary-disk use up sharply on one step | the move is real and is doing the work | | Results unchanged throughout | correctness was never at risk, only cost | The last row is the trap. Because nothing is wrong with the numbers, no data-quality check fires, and the regression is usually noticed as a cost or scheduling problem long after the write that caused it. ## Detecting it on purpose If a job's cost model depends on the move not happening, the absence of the move is part of its contract and should be checked like any other contract: - **Assert at write time** that each output piece contains only the keys the rule assigns to it, on a sample if a full check is too expensive. This catches the appending writer at the moment it offends, not six months later. - **Watch the plan, not the clock.** A check that fails when the redistribution reappears before the join is precise; a runtime threshold is noisy and late. - **Record the agreement where writers will see it** - the key, the rule and the count belong with the table, not in the head of the engineer who set it up. - **Name an owner for writes.** The most durable version of this arrangement is one where a single job is allowed to write the table. ## Deciding what to do once it has drifted Three honest options, and the choice is a cost question rather than a technical one: 1. **Rewrite the offending data** into the agreed layout, and add the guard that stops the next one. Correct, and it costs a full pass over what drifted. 2. **Rewrite both sides at a new agreed count** if the sizes have changed enough that the original count is wrong anyway. More expensive, and the right moment to revisit whether the arrangement still pays. 3. **Let it go.** If the join runs once a week and the move costs minutes, the discipline of policing every writer may cost more than the redistribution does. An arrangement nobody maintains is worse than no arrangement, because the cost model quietly assumes a saving that is gone. The senior instinct being tested here is that this decays by default, so a design that depends on it needs an owner and a check, not just a correct first write.
- Why is this regression usually found late?Because it changes nothing observable about the output. Row counts, values and schemas are unaffected, so quality checks stay green; the only evidence is a redistribution in the plan and a slowly rising runtime and bill, which are easy to attribute to data growth. Detection has to target the plan or the placement, not the results.
- A backfill has to append three months of history. How would you keep the arrangement intact?Have the backfill write through the same path as the regular producer, applying the same key, rule and count, and verify placement on a sample of the output before publishing it. If the backfill cannot do that, treat the table as drifted and plan the rewrite deliberately rather than discovering it later.
- Is it worth keeping when only one of several consumers benefits?Rarely. Every writer pays for the arrangement on every load, while only the join on that one key is repaid. If that join is frequent and large, the arithmetic can work; if it is weekly and small, the policing effort and the frozen count usually cost more than the movement it avoids.
saying these in an interview costs you the question
- Assumes a layout set once stays true forever
- Expects drift to surface as an error or wrong results
- Blames data growth for a runtime that doubled structurally
- Thinks appending records cannot disturb the arrangement
- Proposes a rewrite without adding a guard against the next writer
- Treats a runtime threshold as a reliable detector