How does an MPP planner choose between broadcasting one join input and hash-redistributing both?
answer
- one side copied everywhere, or both re-sorted
- copy cost multiplies by the number of nodes
- the big side never moves in one of them
- every node must hold the copy in memory
- a wrong size estimate is what breaks it
basics
~20 sIt compares network cost. Broadcasting sends the smaller input's bytes to every node, so cost scales with node count; redistributing sends each row of both inputs once. Broadcast wins when one side is small enough to fit in every node's memory.
solid answer
~50 sBoth strategies exist to satisfy the same rule: matching join keys must end up on one node. **Broadcast** sends one input in full to every node — network cost is roughly `size(small side) × node_count`, and every node must hold that copy in memory for its hash table. **Hash redistribution** re-partitions both inputs on the join key — network cost is roughly `size(A) + size(B)` moved once, but nothing needs to be replicated. So the planner estimates the smaller side's cardinality and picks broadcast when `small × N` is well under the redistribution volume and fits in the per-node memory budget. If one input is *already* distributed on the join key, only the other side needs redistributing, which is cheaper still. The whole decision rests on cardinality estimates, so a badly underestimated "small" side is the classic way a query broadcasts gigabytes and spills or fails.
code
text · 12 lines-- broadcast plan: small side replicated, big side stays put
HashJoin ON f.product_id = d.product_id
Scan facts (est 4,000,000,000 rows, local)
Exchange BROADCAST (est 12,000 rows)
Scan dims
-- redistribution plan: both sides re-partitioned on the key
HashJoin ON a.user_id = b.user_id
Exchange HASH(user_id)
Scan events (est 3,000,000,000 rows)
Exchange HASH(user_id)
Scan sessions (est 900,000,000 rows)go deeper
Recall that a distributed join either copies one input to every machine or re-sorts both inputs by the join key, and that the choice depends on how big the smaller input is.
State the two cost formulas — small side times node count versus both sides once — and explain that the broadcast copy has to fit in each node's memory.
Diagnose the failure: a broadcast chosen on a stale or underivable estimate that explodes at runtime. Show how filtering, projecting and refreshing statistics change the decision.
Reason about the interaction with cluster sizing and concurrency: broadcast cost scales with node count and multiplies memory pressure fleet-wide, so plan strategy and capacity decisions are not independent.
## The same goal, two ways to reach it A distributed join needs rows with equal join keys on the same machine. There are exactly two general ways to achieve that with an exchange: 1. **Replicate one side.** Send every row of the smaller input to every node. Each node now holds the whole small side plus its own share of the large side, and can join locally. 2. **Repartition both sides.** Hash both inputs on the join key and send each row to the node owning its bucket. Each node then holds matching subsets of both sides. A third, cheapest case exists when the data is already placed correctly: if both inputs are distributed on the join key, no movement is required at all, and if only one is, only the other side needs to move (a single-sided redistribution). ## The cost model Let `N` be the node count, `S` the bytes of the smaller input after filtering and projection, and `L` the bytes of the larger input. - **Broadcast** moves roughly `S × N` bytes in total (or `S × (N-1)` from each node's perspective, depending on how the engine counts). It is *independent of `L`* — the big table never moves. Every node must materialize a hash table over the whole of `S`, so per-node memory is `S`, not `S / N`. - **Hash redistribution** moves roughly `S + L` bytes, each row exactly once. Per-node memory for the build side is about `S / N`, because the small side is also split. Broadcast therefore wins when `S × N < S + L`, i.e. roughly when `S` is small relative to `L / N`. Concretely: with 20 nodes, broadcasting a 2 GB input moves 40 GB, while redistributing a 2 GB and a 4 TB input moves about 4 TB — broadcast is dramatically cheaper. Broadcasting a 200 GB input across the same 20 nodes moves 4 TB and needs 200 GB of hash table on every node, which is usually fatal. Note the asymmetry with cluster size: **broadcast cost grows with node count, redistribution cost does not.** A broadcast plan that is comfortable on a 4-node cluster can become the dominant cost on a 64-node one, which is one reason the "same query, bigger cluster, no faster" complaint appears. ## What the planner actually uses The decision is made from estimated cardinalities and average row widths, after predicate and projection pushdown — what matters is the size of the input *as it arrives at the join*, not the table's size on disk. A highly selective filter can turn a large table into a broadcastable input. Most engines also apply a hard ceiling: if the estimated broadcast side exceeds a memory or size threshold, redistribution is chosen regardless. Because the choice is estimate-driven, its failure mode is estimate-driven too. If statistics are stale, or the input is the result of an earlier join or an aggregation whose cardinality is hard to predict, the planner may believe a side is 50 MB when it is 50 GB. The query then broadcasts 50 GB to every node, memory pressure spikes on all of them at once, the join spills or the query is killed for exceeding its memory limit. The tell in a profile is a broadcast exchange whose actual row count dwarfs its estimate. ## Row width matters as much as row count Shuffle cost is bytes, not rows. Selecting fifteen columns where three are needed multiplies the broadcast volume accordingly, and wide string columns dominate quickly. Projecting only the join key and the columns actually consumed above the join — and filtering before the exchange rather than after — often changes the decision from redistribution to broadcast and shrinks the moved data at the same time. ## Skew interacts with the choice Redistribution is vulnerable to skew: if one join key value accounts for a large share of rows, one node receives that whole share and becomes the query's runtime. Broadcast is immune to key skew on the probe side, because the large input never moves — each node simply probes its own local rows against the full copy. That makes broadcast attractive when the join key is known to be lopsided and the other input is genuinely small. Conversely, broadcast has its own skew-like hazard: it multiplies memory pressure uniformly, so it fails on *all* nodes at once rather than on one, and it competes with every other query running concurrently on the same cluster. ## Practical levers An engineer influences this in four ways: keep statistics current so the estimate is right; filter and project before the join so the candidate side really is small; align physical placement with the hottest join key so neither side has to move; and, where the engine exposes it, review the chosen strategy in the plan and treat an unexpected broadcast of a large input as a statistics bug rather than a tuning knob.
- Why can a broadcast join that works well on a small cluster become the bottleneck on a much larger one?Broadcast cost is the small side's size multiplied by the node count, so it grows linearly as you add nodes, while redistribution cost stays roughly constant. Doubling the cluster halves the scan time but doubles the broadcast traffic and duplicates the same hash table on twice as many nodes, so past a point the added nodes cost more in network and memory than they return in parallelism.
- What happens when the planner badly underestimates the size of the input it chose to broadcast?Every node tries to build a hash table over far more data than budgeted. The join spills to disk on all nodes simultaneously, or the query is killed for exceeding its memory limit, and the network carries the oversized payload N times over. The fix is refreshing statistics or rewriting so the estimate is derivable, not raising the memory limit.
Broadcasting is mailing a photocopy of a thin reference booklet to every branch office; redistribution is shipping both filing cabinets to a central sorting depot so matching folders meet. The booklet is cheap to copy; the cabinets are not.
saying these in an interview costs you the question
- Thinking broadcast is always cheaper because one side is small
- Ignoring that every node must hold the broadcast copy in memory
- Judging size by table size on disk, not by rows after filtering
- Assuming the choice is independent of cluster size
- Treating a bad broadcast as a memory-limit problem, not a statistics problem