A star join on a 6 TB fact table runs 40x slower after a dimension grew. How do you diagnose it?
answer
- read the plan that ran, not the one predicted
- is every node hurting, or just one?
- estimated rows versus actual rows on the build side
- a placeholder key can own half the fact table
- check how much of the fact table is still being skipped
basics
~20 sCompare estimated versus actual rows on the join's build side and check whether the plan still broadcasts. A dimension that outgrew the broadcast threshold either replicates far too much data to every node or flips to redistribution, where a skewed join key can pin the work onto one node.
solid answer
~50 sStart from the plan, not the query. Two distinct failures produce this symptom. First, the dimension is still broadcast but is now large: every node receives the full copy, per-node memory spikes, the build side spills, and the pain is **uniform across all nodes**. Second, the planner has flipped to hash redistribution: now both sides move, and if the join key is skewed one node receives a disproportionate share and becomes a straggler — the pain is **concentrated on one node** while others finish early. Distinguish them by per-node time and memory in the profile. Then check estimates: if the estimated build-side rows are orders of magnitude below actual, stale statistics on the dimension are the root cause and refreshing them may restore a sane plan. Finally check whether a runtime filter is still being produced; a dimension filter that lost selectivity as the table grew silently stops pruning the fact scan.
code
text · 5 lines-- per-node profile excerpt, hash-redistributed join
node 03 rows 1,940,000,000 time 612s peak-mem 41 GB spill 18 GB
node 04 rows 52,100,000 time 19s peak-mem 3 GB spill 0
node 05 rows 49,800,000 time 18s peak-mem 3 GB spill 0
node 06 rows 51,300,000 time 18s peak-mem 3 GB spill 0go deeper
Recall the first move: look at the query plan that actually ran and compare how many rows each step expected against how many it produced.
Explain the two failure shapes — a replicated dimension that got too big versus both sides being re-partitioned on a lopsided key — and why the second concentrates on one node.
Demonstrate the full diagnosis: per-node timing and memory to classify the failure, estimate-versus-actual to find stale statistics, rows-read-to-rows-kept to see whether pruning died, then fixes ordered from cheapest to structural.
Own the prevention side. Star-join plans flip at a size threshold exactly once and silently, so decide what plan-shape and bytes-scanned monitoring exists, and which tables get layout investment before they cross that line.
## Why this question is asked It is the archetypal warehouse incident: nothing in the SQL changed, data volumes grew modestly, and a query fell off a cliff. Star joins are unusually cliff-prone because their plans depend on a size estimate for the dimension, and there are threshold effects on both sides of that estimate. ## Step 1 — get a plan with actual counters, not just estimates An estimated plan tells you what the optimizer believed; you need what happened. Whatever the engine's profiling facility is, you want per-operator **actual rows**, **bytes exchanged**, **peak memory**, **spill volume**, and ideally **per-node timing**. Almost every diagnosis below is a comparison between two of these numbers. ## Step 2 — classify the failure by its blast radius The single most useful discriminator is whether the slowdown is spread evenly or concentrated. **Uniform across nodes → broadcast blow-up.** The dimension is still being replicated, but it is now big. Every node allocates the same oversized hash table, every node hits its memory limit, every node spills to disk at roughly the same time. Network bytes are `size(dimension) × node_count` and have grown linearly with the dimension. Memory graphs look like every node doing the same painful thing simultaneously. **Concentrated on one or two nodes → skew in a redistributed join.** The planner decided the dimension is no longer broadcast-worthy and switched to hashing both sides on the join key. Now the distribution of key values matters. If 40% of fact rows carry one `customer_id` (a default, an "unknown" placeholder, a house account, or a NULL surrogate mapped to a single key), the node owning that hash bucket does 40% of the join work alone. Cluster utilization looks terrible: one node pinned, the rest idle, and total runtime set by the slowest. The "unknown member" key is worth calling out specifically. Star schemas conventionally map missing foreign keys to a single sentinel dimension row. That sentinel can accumulate an enormous share of the fact table, and it is invisible until the join stops being a broadcast. ## Step 3 — compare estimated to actual on the build side If the estimate says 30,000 rows and the actual is 90 million, the planner was not making a bad decision — it was making a good decision from bad inputs. Common causes: - Statistics on the dimension have not been refreshed since it grew. - A predicate on the dimension is estimated with an independence assumption that does not hold, so the filter's true selectivity is far worse than modelled. - The dimension is now populated by a pipeline that lands data in a shape statistics do not describe well. Refreshing statistics is the cheap first move and frequently the whole fix. ## Step 4 — check whether pruning still happens A star join's fact scan is usually kept small by a runtime filter derived from the dimension's surviving keys. Two things break it as a dimension grows: - The filter passes many more keys, so its selectivity collapses and block skipping stops. - The engine adaptively disables the filter after observing a poor pass rate. The tell is the ratio of fact rows read to rows surviving the join. If that ratio was 40:1 and is now 1.2:1, the fact table is being fully scanned where it used to be pruned, and no amount of join-strategy tuning fixes it. The remedy is on the layout side — clustering or partitioning the fact table on the key that dimension filters actually target. ## Step 5 — decide the fix, in order of cheapness 1. **Refresh statistics.** Free, fast, often sufficient. 2. **Make the dimension filter selective again, or filter it earlier.** If a subquery or view is materializing the whole dimension before filtering, restructure so the predicate applies before the exchange. 3. **Handle the skew explicitly.** If a sentinel key dominates, split the query: process the sentinel rows separately (they usually need no dimension attributes at all) and union the result. Salting the key is the heavier alternative when a genuine business key is hot. 4. **Force the strategy.** If you know the dimension is genuinely small after filtering and the planner does not, pushing it back to a broadcast can be right — but treat a forced strategy as a pin that will itself go stale. 5. **Change the layout.** Sort/cluster the fact table on the filtered key, or co-locate the pair if the engine supports it and this join dominates the workload. 6. **Reconsider the shape.** If this dimension has stopped being small permanently, flattening its few needed attributes into the fact table removes the join entirely. ## Step 6 — prevent the recurrence The underlying lesson is that star-join plans are threshold-sensitive. A monitoring signal worth having is the ratio of bytes scanned to rows returned per recurring query, and an alert on plan-shape changes for the queries that matter. A dimension that is growing steadily will cross the broadcast threshold exactly once, quietly, and the first symptom will be a production incident. ## Interview framing Lead with the discriminator — uniform pain means broadcast, concentrated pain means skew — because it shows you have actually read a profile rather than memorized a list of causes. Then walk estimates, then pruning, then fixes ordered by cost. Mentioning the sentinel/unknown-member key as a skew source is the detail that marks someone who has operated a real star schema.
- How do you tell a broadcast blow-up from key skew using only per-node metrics?Broadcast blow-up is uniform: every node allocates the same oversized build-side hash table, so memory peaks and spill volume look alike everywhere and all nodes slow together. Skew is concentrated: one node's row count, runtime and memory dwarf the rest while the others finish and idle. Total runtime tracking the single slowest node is the signature of skew.
- Why is the star schema's unknown-member key a classic skew source?Dimensional modelling maps missing or unresolved foreign keys to one sentinel dimension row rather than leaving NULLs. Over time a large share of fact rows can carry that single key. While the dimension is broadcast this is invisible, because no redistribution happens. The moment the join flips to hash redistribution, every one of those rows hashes to the same node.
- What would you change if the fix has to survive further dimension growth?Stop depending on the dimension staying broadcast-small. Either flatten the handful of attributes the query needs into the fact table so the join disappears, or change the fact table's layout so the filtered key prunes blocks directly, or split the sentinel-key traffic out of the main join. Forcing a broadcast is a pin that will break again at the next growth step.
- When is refreshing statistics not enough even though estimates were wrong?When the true situation is genuinely bad. Accurate statistics only guarantee the planner sees reality; if reality is a 90-million-row build side or a key where one value holds 40% of the fact table, the best available plan is still slow. At that point the fix is structural — layout, skew handling, or removing the join — not statistical.
saying these in an interview costs you the question
- Jumps to adding an index without reading the plan
- Assumes any slow join means the cluster needs more nodes
- Cannot distinguish uniform memory pressure from single-node skew
- Thinks refreshing statistics fixes skew
- Ignores whether the fact scan is still being pruned at all