You spread a heavy join key on one input with an artificial suffix - what must the other input do, and what does that cost?
answer
- both sides must still agree on the key
- spread one side, expand the other
- copies multiply the expanded input
- the partial results are disjoint
- outer joins duplicate the expanded side
basics
~20 sEvery row on the other input under a spread key must be copied once per suffix value, so each pairing can still meet. That multiplies those rows by the width, and the combining pass is a union rather than a re-aggregation.
solid answer
~50 sAn equality join sends both inputs' rows for a key to the same destination, so changing one side's key to `(key, s)` leaves the other side's plain key matching nothing. The invariant is: **spread one side, expand the other**. Each row on the spread side gets exactly one suffix, random or round-robin; each row on the other side whose key is being spread is replicated to all `W` suffix values, one copy each. Then `(key, s)` meets `(key, s)`, and every original pair is produced exactly once. The second pass is therefore not a fold. The `W` partial join outputs are disjoint, so you drop the suffix and concatenate them - no combining function, no equivalence question. The cost is the replication: the expanded side's rows for those keys are multiplied by `W`, which is the reason to spread only the keys measured heavy.
go deeper
Remember that a join matches on equal keys, so changing the key on one side alone loses matches unless the other side is given every suffix value in turn.
Explain the invariant - one suffix per row on the spread side, all suffixes on the other - and why that makes the partial results disjoint and the final step a union instead of a fold.
Bound the replication with a measured list of heavy values, and work out the outer-join case explicitly: the expanded side is the one that duplicates unmatched rows.
Judge whether a join carrying a hand-written spread and a maintained heavy-key list belongs in a platform at all, against reshaping the inputs or accepting the slower run, and record the assumption it rests on.
## Why a join needs more than an aggregation does A join on equality cannot be computed from what one worker already holds, so it opens with a **redistributing step**: both inputs' rows are sent to the destination that owns their key, chosen by a function of the key value alone, and the pairing happens where they meet. A **heavy key** - one value covering a large share of one or both inputs - therefore concentrates on one destination, and that destination sets the step's duration. The remedy is the same rewrite as for a grouping, with one extra obligation. Appending a suffix `s` in `0..W-1` to the key of one input scatters that input's rows, but it also changes the value that decides matching. A row on the other input still carrying the plain key now agrees with nothing, and the join quietly returns fewer rows than it should. ## The invariant **One side is spread; the other is expanded.** - On the spread side, each row of a heavy key receives exactly one suffix, drawn at random or handed out round-robin. - On the other side, each row whose key is being spread is replicated into `W` copies, carrying `s = 0, 1, ... W-1` in turn. Now consider any pair that the original join would have produced, between a row `L` on the spread side and a row `R` on the other. `L` carries one suffix `s`; `R` exists under every suffix, including `s`; the two composite keys are equal, so they meet at that destination and the pair is produced. It is produced only there, because `L` exists nowhere else. So the `W` partial joins are disjoint and their union is exactly the original result. This is why the second pass is a concatenation rather than a fold: nothing needs combining, only the suffix needs dropping. The correctness question that dominates a spread aggregation - which aggregates survive two passes - does not arise for an inner join at all. ## What it costs | | spread side | expanded side | |---|---|---| | rows under a spread key | unchanged | multiplied by `W` | | bytes moved for those keys | unchanged | multiplied by `W` | | destinations reached | up to `W` | up to `W` | | output rows | unchanged | unchanged | The output is the same size as before; the input to the join is not. Expanding a whole input by a width in the hundreds is usually far worse than the imbalance, so the replication is applied only to rows whose key appears on a measured list of heavy values, and everything else passes through with a constant suffix. That keeps the multiplication proportional to the few keys that caused the problem. A second limit: if both inputs are heavy on the same value, the pairing itself is multiplying rows, and no key rewrite changes how many pairs exist. The rewrite spreads the work of producing them; it does not make fewer of them. ## Outer joins are where this goes wrong The invariant is asymmetric, so outer semantics are too: - **Left outer, with the left side spread.** Each left row exists exactly once, under one suffix. If it matches, it matches there; if it does not, it is emitted once, null-extended. The semantics survive unchanged. - **Right outer, with the right side expanded.** Each right row now exists in `W` copies. A copy that lands on a suffix where no left row for that key happens to sit is unmatched and gets null-extended, so a right row with no partner can be emitted up to `W` times, and one with a partner can be emitted both matched and null-extended. The result is wrong in a way no downstream count will explain. - **Full outer** inherits the right-side problem. The repairs are to compute the null-extended side separately - identify unmatched keys on the expanded side once, outside the spread path, and union them in - or to choose which side to spread so that the preserved side is the one spread rather than the one expanded. ## What varies between runtimes Do not assume the runtime will do any of this. Some engines re-decide the part of a plan they have not run yet from the sizes a finished step actually produced, and will divide an oversized join share by themselves - for some join types only, and not for every one. Others never rewrite a plan at all. Over an endless input, where no step finishes and there is nothing completed to measure, this kind of automatic help is largely absent, and a join there carries the further constraint that the rows being expanded must be retained somewhere for as long as their partners may arrive. Write the rewrite against what your runtime actually measures, and state that assumption where the next reader will find it.
- Why not attach an independently drawn suffix to both sides?Because matching needs the two composite keys to be equal, and two independent draws agree only by chance - roughly one pairing in `W` survives. The asymmetry is the point: one row on the spread side must exist under exactly one suffix, and its partners must exist under all of them.
- Why is the second pass a union rather than a fold, as it is for an aggregation?Because the `W` partial joins are disjoint. Every original pair is produced in exactly one of them, since the spread side's row lives under one suffix only. Dropping the suffix and concatenating the partials reproduces the original result, so no combining function and no associativity question is involved.
- How do you keep the replication from dominating the job?Replicate only the rows whose key is on a measured list of heavy values, and let everything else through with a constant suffix. The expansion is then bounded by `rows under the heavy keys x W` rather than by the whole input, which is usually a small fraction of it.
saying these in an interview costs you the question
- Give both sides a random suffix and join on the composite key
- Replicating the other side does not change how much data moves
- The partial join outputs have to be re-aggregated like a sum
- An outer join behaves the same as an inner one under this rewrite
- Spreading the key reduces how many output rows the join produces