skip to content

A report query that returned in milliseconds during testing now runs for hours in production on the same plan shape — an index-driven nested loop join. What property of that operator makes it so sensitive, and how would you reason about the regression?

level: seniorimportance: must knowfreq 50%

answer

  1. cost = outer rows x per-probe constant
  2. the estimate is the load-bearing assumption
  3. no fallback mid-execution
  4. estimated vs actual on the OUTER branch
  5. correlated predicates underestimate

basics

~20 s

Its cost is the outer row count times a per-lookup constant, and the optimizer chose it believing the outer side would produce few rows. If the real outer count is far larger, cost grows linearly with no fallback: millions of random index probes instead of one sequential scan plus a hash join.

solid answer

~60 s

An index nested loop is only cheap because the outer side is small. The optimizer estimated the outer cardinality, multiplied it by a per-probe cost, and won the comparison against hash and merge join. That estimate is the single load-bearing assumption, and the operator has no fallback when it is wrong: at a thousand times the expected outer rows you get a thousand times the probes, each a scattered read that no longer hits cache. My first check is estimated versus actual rows on the outer side, and how many times the inner side executed with how many rows per execution. Common causes: production data volume or skew unlike the test set, stale or missing statistics, correlated predicates the optimizer treats as independent, a parameter value far from the typical one, and an inner index that is no longer usable. Fixes range from refreshing statistics and adding a covering index to reducing the outer side with a better predicate, or steering the plan to a hash join when the join really is large-to-large.

go deeper

for a junior

Say the loop repeats once per outer row, so many more outer rows than expected means proportionally more work.

for a middle

Add that the optimizer chose it based on an estimated small outer cardinality, and that you would compare estimated with actual rows.

for a senior

Walk the diagnosis end to end — actual plan, loop executions, rows per execution, statistics, correlation, cache differences — and pick the fix that matches the cause.

for a principal

Discuss the crossover regime, robustness of plan choice under data growth, testing on production-shaped data, and when plan stability is worth its cost.

## Why this operator is the fragile one Every join algorithm has a cost model, but their sensitivities differ. A hash join reads each input once; if the estimate is off, cost is wrong by a constant factor and possibly some spilling. A merge join over sorted inputs is likewise a linear pass. An index nested loop's cost is n_outer multiplied by a per-probe constant, and that constant is not small: an index descent plus, usually, a random base-table fetch per matching row. Being wrong about n_outer scales the whole join, and unlike the other operators there is no point at which the strategy switches to something more suitable — execution follows the plan that was chosen. ## The crossover it sits near There are two regimes. Few outer rows: probing is dramatically cheaper than reading the inner table, and the loop wins by orders of magnitude. Many outer rows: the probes collectively touch much of the inner table anyway, in random order, and one sequential scan plus a hash build would have been far cheaper. The cardinality estimate places the query on one side of that crossover. Small errors near the crossover cost a little; large errors cost hours. ## What actually differs between test and production Volume is the obvious one: a test dataset of a few thousand rows keeps everything in cache, so even a bad plan looks instant. Distribution matters more — predicates that are selective on synthetic data may match a large fraction of real rows, and skew means the average case in statistics is not the case being run. Correlated predicates are a classic trap: optimizers typically multiply selectivities as if conditions were independent, so two correlated filters produce an estimate far below reality. Parameter values matter when a plan chosen for one value is reused for another with very different selectivity. Statistics may be stale after bulk loading. And the inner index may be unusable in production because of a type mismatch, an expression wrapped around the column, or an index that was never created there. ## How to diagnose Get an execution plan with actual row counts, not just estimates. Compare estimated with actual on the outer branch first — that is the assumption under test. Then look at how many times the inner branch executed and the rows returned per execution; multiplied together they give the real work done. If actuals dwarf estimates on the outer side, the diagnosis is a cardinality problem, not an indexing problem. If executions are as expected but each is slow, the inner access path or cache behaviour is at fault instead. Comparing buffer or read counts between environments separates cache effects from genuine algorithmic cost. ## Fixes, roughly in order Refresh statistics and, where supported, add multi-column or expression statistics so correlated predicates estimate better. Reduce the outer side: a more selective predicate, a filter applied earlier, or a rewrite that makes the small side genuinely small. Make the inner probe cheap with a covering index so the random base-table fetch disappears. If the join truly is large-to-large, the fix is a different operator — make hash or merge join viable and, if necessary, steer the optimizer there. Plan stability mechanisms are a last resort: they freeze today's choice, which is only right if you are confident the data shape will not move. ## The prevention framing The deeper lesson is that plans chosen on the assumption of a small outer input are a latent risk on any table that grows. Testing on production-shaped data volumes and distributions, not just production schema, is what surfaces this before users do. Monitoring queries by total work rather than average latency catches the skewed parameter values that only occasionally take the expensive path.

  • Why do correlated predicates so often lead to this failure?
    Optimizers commonly assume predicate independence and multiply selectivities, so two conditions that in reality select the same rows produce an estimate far below the truth. That underestimated outer cardinality is exactly what makes an index nested loop look cheap, and the plan then does far more probes than costed. Multi-column or expression statistics, where available, correct the estimate.
  • How does an index nested loop's failure mode differ from a hash join's?
    A hash join with a bad estimate still reads each input once; the penalty is a wrong memory grant leading to partitioning and spilling, which is bad but bounded. An index nested loop with a bad estimate does proportionally more random probes with no bound short of finishing, so its degradation is unbounded and far steeper.

saying these in an interview costs you the question

  • Jumping straight to adding indexes without checking estimated versus actual rows on the outer side
  • Blaming the database engine rather than the cardinality assumption the plan rests on
  • Assuming the same plan shape implies the same cost — the loop count is what changed
  • Pinning or hinting the plan as a first move instead of fixing the estimate
  • Testing on production schema with toy data volumes and calling that verification

context