skip to content

When joining a huge fact table to a small dimension, why does an MPP engine broadcast the dimension?

level: middleimportance: must knowfreq 72%

answer

  1. only one side has to move
  2. the fact table is the expensive side
  3. copies of a small table are cheap
  4. the multiplier is the node count
  5. estimates decide which plan you get

basics

~20 s

Broadcasting sends one full copy of the small dimension to every node so each node joins its local fact rows without moving them. The fact table is orders of magnitude larger, so shipping the dimension costs far less network traffic than redistributing both sides.

solid answer

~50 s

In a shared-nothing MPP engine, an equi-join can only be computed where matching rows meet on the same node. There are two ways to arrange that: **redistribute** both inputs by a hash of the join key, or **broadcast** the smaller input in full to every node. A fact table may hold billions of rows spread across all nodes; a dimension may hold thousands. Broadcasting ships `dimension_size × node_count` bytes and leaves the fact table completely still — no shuffle of the expensive side at all. Each node then builds a local hash table on its copy of the dimension and probes it with its local fact rows, which is also the ideal shape for a pipelined hash join. Redistributing both sides would move the entire fact table across the network for no benefit. The decision hinges on the planner's size estimate for the dimension, and it falls apart when that estimate is badly wrong.

go deeper

for a junior

Recall the basic picture: data lives on many machines, and rows can only be joined when they sit together. Sending a copy of a tiny lookup table everywhere is cheaper than moving a billion-row table.

for a middle

Explain both physical shapes and their cost formulas: broadcast moves size(small) × nodes, redistribution moves both inputs in full. Say why the small side becomes the hash-table build side and the fact table the probe side.

for a senior

Show you have watched a broadcast go wrong. Talk about estimate-driven plan choice, per-node memory limits, the uniform-across-all-nodes signature of a broadcast blow-up, and how you would force or avoid the strategy.

for a principal

Own the trade-off at cluster scale: as node counts rise the broadcast multiplier grows, co-location can only be spent on one join per table, and the policy question is which joins get physical layout investment versus which are left to the planner.

## The problem an MPP join has to solve An analytical warehouse spreads each table's rows across many nodes (or slices/workers). For an equi-join `fact.dim_id = dim.id`, two matching rows can only be joined if they are physically on the same node at the same time. Unless the two tables happen to already be co-located on the join key, the engine must move data. The plan expresses this with an *exchange* operator, and there are essentially two shapes of exchange for a join. ## Shape 1 — hash redistribution Both inputs are re-partitioned by `hash(join_key) % node_count`. Every row of both tables travels over the network to the node that owns its key bucket. This is symmetric and always correct, but its cost is proportional to the size of *both* inputs. For a star join that means moving the whole fact table. ## Shape 2 — broadcast (replicated join) One input — the small one — is sent in its entirety to *every* node. The large input never moves. Network cost is `size(small) × N` where `N` is the number of nodes; the large side contributes zero network bytes. Each node then runs an independent local join: ```sql SELECT d.country, SUM(f.amount) FROM fact_sales f JOIN dim_customer d ON f.customer_id = d.customer_id GROUP BY d.country; ``` If `dim_customer` is 50 MB and `fact_sales` is 4 TB, broadcasting 50 MB to each of 32 nodes moves ~1.6 GB. Redistributing would move ~4 TB. The asymmetry is why broadcast is the default plan for a star join and why the pattern is sometimes called a *replicated* or *map-side* join. ## Why it pairs so well with a hash join After the broadcast, each node has a complete copy of the dimension — small enough to build an in-memory hash table keyed on the dimension's key. The node then streams its local fact partition through that hash table as the probe side. Nothing blocks except the (tiny) build; the fact scan stays pipelined straight into aggregation. That is the cheapest physical shape available for this query, and it also enables runtime filtering: because the build side is materialized first, the engine can derive a filter from the dimension keys and push it down into the fact scan. ## Where broadcast goes wrong The whole argument rests on "the dimension is small." Three failure modes: 1. **Bad estimates.** The planner estimates the build side from statistics. If a filter on the dimension is estimated at 1,000 rows but actually yields 40 million, the engine broadcasts 40 million rows to every node — `N` times the memory and `N` times the network. This is the classic MPP blow-up: memory pressure on every node at once, spilling, or an out-of-memory abort. 2. **Broadcast is per-node memory, not cluster memory.** A 2 GB dimension broadcast to 32 nodes consumes 2 GB *on each node*, not 2 GB total. The threshold for "small" is set by per-node memory, not cluster memory. 3. **High node counts.** As `N` grows, `size(small) × N` grows linearly. On a very wide cluster, a moderately sized dimension stops being cheap to broadcast. ## When redistribution is the right call When both sides are large — a fact-to-fact join, or a fact joined to a dimension that is itself hundreds of millions of rows — redistribution wins because broadcast's `× N` multiplier dominates. The engine hashes both sides on the join key, and each node handles one slice of the key space. The risk shifts from memory blow-up to **skew**: if one key value accounts for a large share of rows, one node receives a disproportionate share of the work and becomes a straggler while the rest idle. ## The third option: co-location Some engines let you physically arrange two tables so that rows with the same join key already live on the same node. Then neither side moves and the join is entirely local. This is powerful for a repeatedly joined pair, but you only get one such arrangement per table, so it is spent on the most important join and everything else still broadcasts or redistributes. ## What to say in an interview Name the three physical strategies — co-located, broadcast, redistribute — say that cost is dominated by bytes crossing the network, and explain that broadcast trades `N` copies of the small side for zero movement of the big side. Then show you know the failure mode: broadcast is only safe while the estimate for the small side is right, and a wrong estimate hurts every node simultaneously.

  • How would you recognize from a query profile that a broadcast blew up because the small side was under-estimated?
    Look for a huge gap between estimated and actual rows on the broadcast/build operator, then for the symptoms that follow: high memory on every node at once rather than one, spill-to-disk or a memory-limit abort, and network bytes far above what the dimension's true size predicts. Uniform pain across all nodes points at broadcast; pain on one node points at skew in a redistributed join.
  • Why does broadcasting stop being attractive as the cluster gets wider?
    Broadcast network and memory cost scale as size(small) × node_count, while redistribution cost stays proportional to the total data size regardless of node count. Doubling the cluster doubles the broadcast bill for the same dimension. On very wide clusters a dimension that was cheap to replicate on eight nodes becomes expensive on a hundred, and the planner's crossover point shifts toward redistribution.
  • If both sides of a join are large, what replaces broadcast and what new risk appears?
    Hash redistribution: both inputs are re-partitioned by a hash of the join key so matching rows land on the same node. Memory blow-up is no longer the main risk — skew is. If one key value dominates, one node receives most of the rows, runs far longer than the rest, and becomes a straggler that determines the whole query's runtime while other nodes sit idle.

saying these in an interview costs you the question

  • Says the engine always broadcasts the second table listed
  • Thinks broadcast memory is shared across the cluster, not per node
  • Claims broadcast is always faster because it avoids the network
  • Believes join strategy is chosen from table row counts on disk only
  • Says redistribution is obsolete in modern warehouses

context