skip to content

questions

5

A single query joins eight tables. Why does the order in which the database engine joins them matter so much, and roughly how many possible orderings exist?

level: middleimportance: must knowfreq 55%

answer

  1. Same rows, different intermediate sizes
  2. n! left-deep, (2n-2)!/(n-1)! bushy
  3. 8 tables = 40,320 vs ~17M
  4. Push selectivity down early
  5. Search, not enumerate

basics

~20 s

Join order never changes the result, but it changes the size of the intermediate results carried between operators, so cost can differ by orders of magnitude. The number of orderings grows factorially: eight tables give 40,320 left-deep orders and roughly 17 million bushy ones.

solid answer

~50 s

Joins are logically associative and commutative, so every ordering returns the same rows. What changes is the **size of intermediate results**. Joining the two most selective tables first can leave a few hundred rows to carry through the rest of the plan; starting with two large unfiltered tables can materialize hundreds of millions of rows that later get discarded. That is easily a 1000x difference in CPU, memory and I/O. The search space is factorial. Restricted to left-deep shapes it is `n!` — 40,320 for eight tables, ~479 million for twelve. Allowing bushy shapes it is `(2n-2)!/(n-1)!` — about 17 million for eight. Each join node then multiplies that again by its physical operator choice (nested loop, hash, merge) and each table by its access path. So an optimizer cannot enumerate everything. It uses dynamic programming with pruning, avoids cross products, and switches to heuristics past a join-count threshold.

code

sql · 5 lines
sql
SELECT oi.sku, oi.qty
FROM   customer c
JOIN   orders      o  ON o.customer_id = c.id
JOIN   order_items oi ON oi.order_id   = o.id
WHERE  c.email = '[email protected]';

go deeper

for a junior

Know that join order doesn't change results but hugely changes speed, and that the engine — not the FROM clause order — decides.

for a middle

Give the factorial growth, quote n! for left-deep versus the much bigger bushy count, and explain intermediate cardinality as the cost driver.

for a senior

Connect the space size to why optimizers use dynamic programming, prune cross products, and fall back to heuristics; note planning time as a real cost.

for a principal

Frame it as a search-versus-guarantee tradeoff: bounded planning budget buys good-not-optimal plans, which is why very wide queries need decomposition or plan reuse rather than more search.

## What "join order" means A query mentioning tables A, B, C, D does not say how to combine them. The engine must choose a binary tree of join operators: `((A⋈B)⋈C)⋈D`, or `(A⋈B)⋈(C⋈D)`, or `((D⋈C)⋈A)⋈B`, and so on. Because the relational join is associative and commutative, **all of these produce exactly the same set of rows**. Choosing among them is therefore purely a performance decision — which makes it the single highest-leverage decision a cost-based optimizer makes. ## Why cost varies so wildly A join operator's cost is driven mostly by how many rows arrive at it. Rows arriving are the output of whatever ran below. So the order determines the *cardinality* of every intermediate result. Consider `orders` (100M rows), `order_items` (500M), and `customer` (1M) with a filter selecting one customer. Joining `customer` (filtered to 1 row) to `orders` first yields maybe 50 rows, and joining those 50 to `order_items` yields ~300 rows: the whole query touches a few hundred rows through indexes. Join `orders` to `order_items` first and you materialize a 500-million-row intermediate before the customer filter can throw nearly all of it away. Same answer; hours versus milliseconds. The general heuristic behind good orders is **push selectivity down**: apply filters early and join in an order that keeps intermediate results small, so the reducing power of each predicate is exploited as early as possible. ## Counting the space Two numbers matter. *Left-deep orders* (every join's right input is a base table) are just permutations of the tables: `n!`. | tables | n! (left-deep) | |---|---| | 4 | 24 | | 6 | 720 | | 8 | 40,320 | | 10 | 3,628,800 | | 12 | 479,001,600 | *All shapes (bushy included)* multiply permutations by the number of binary tree shapes, the Catalan number C(n-1), giving `(2n-2)!/(n-1)!`: 120 for four tables, ~17.3 million for eight, ~17.6 billion for ten. And that is only the logical order. Each of the n-1 join nodes independently picks a physical algorithm, and each leaf picks an access path (sequential scan, index scan, index-only scan). The true plan space is the product of all three, so it explodes far faster than the table above suggests. ## How real optimizers cope 1. **Dynamic programming (System R / Selinger).** Build the best plan for every *set* of tables bottom-up, reusing sub-results. This turns a factorial problem into roughly `O(3^n)` time and `O(2^n)` memory — still exponential, but tractable to ~10–15 tables. 2. **Shape restriction.** Many optimizers search only left-deep (or left-deep plus limited bushy) trees, cutting the space from `(2n-2)!/(n-1)!` down to `n!`. 3. **No cross products.** Orders that join two tables with no join predicate between them are pruned unless nothing else is possible. On a chain-shaped join graph this is a massive reduction; on a clique-shaped graph it saves nothing. 4. **Pruning by cost.** Any partial plan already more expensive than a known complete plan for the same table set is discarded. 5. **Heuristic fallback.** Past a configured join count, engines switch to greedy construction or a randomized/genetic search that samples the space instead of covering it. ## The consequence engineers feel Because the space is searched, not enumerated, two things follow in production. First, **planning time itself becomes a cost** on very wide queries — sometimes exceeding execution time. Second, **plan quality degrades gracefully rather than being guaranteed**: past the threshold you get a good plan, not the best plan under the cost model, and it may not be the same plan every time. Both are reasons to keep the number of joined relations in a single statement bounded, and to reuse prepared plans for hot queries.

  • If join order never changes the result, why can't the engine just pick any order and be done?
    Because cost is not order-invariant even though output is. Each ordering produces different intermediate cardinalities, and every operator above pays for those rows in CPU, memory and I/O. Picking arbitrarily risks materializing a multi-hundred-million-row intermediate that a different order would never have built.
  • Which join-graph shapes make the search space cheaper to explore?
    Chain (linear) graphs are cheapest: excluding cross products leaves only a small fraction of the permutations valid. Star graphs are next. Clique graphs — where every table has a predicate to every other — are worst, because no ordering can be pruned as a cross product and the optimizer must consider nearly the full space.
  • Does adding a physical operator choice change the counting?
    Yes, multiplicatively. Each of the n-1 join nodes can be a nested loop, hash or merge join, and each base table has several access paths. The optimizer's real search space is the logical order space times all of those, which is why pruning partial plans by cost matters as much as restricting shapes.

Like planning errands: every route visits the same shops, but doing the one that shrinks your shopping list first means you carry less for the rest of the trip.

saying these in an interview costs you the question

  • Claiming join order changes the result set, so the optimizer must preserve the written order
  • Believing the FROM-clause order is the execution order in a cost-based engine
  • Saying the space is only n! regardless of tree shape — that is the left-deep count only
  • Assuming the optimizer always finds the globally cheapest plan
  • Ignoring that intermediate result size, not table size, drives cost

context

open as a page

What is the difference between a left-deep and a bushy join tree, and why do many query optimizers restrict their search to left-deep plans?

level: middleimportance: should knowfreq 38%

basics

~20 s

In a left-deep tree every join's right input is a base table, so joins form a pipeline. In a bushy tree both inputs can be intermediate results. Left-deep shrinks the search space from about (2n-2)!/(n-1)! to n! and pipelines well, but bushy plans can be far better for star-shaped queries and parallel execution.

open as a page

Many relational optimizers stop doing exhaustive join-order search once a query exceeds a certain number of joined tables and switch to a greedy or randomized/genetic algorithm. Why, and what do you give up?

level: seniorimportance: should knowfreq 30%

basics

~20 s

Exhaustive dynamic programming is exponential — roughly 3^n time and 2^n memory — so past about a dozen tables the planning time and memory exceed any benefit. Engines switch to greedy construction or randomized/genetic search. You give up the guarantee of the cheapest plan under the cost model, and randomized search also gives up plan determinism.

open as a page

Explain how the classic System R (Selinger) dynamic-programming algorithm chooses a join order, and why it is cheaper than trying every ordering.

level: seniorimportance: should knowfreq 40%

basics

~20 s

It builds plans bottom-up over subsets of tables: best plan for every single table, then every pair, then every triple, keeping only the cheapest plan per subset. Reusing sub-results turns a factorial search into roughly 3^n time and 2^n memory. It also keeps extra plans that produce a useful sort order.

open as a page

A nightly reporting statement joins around 25 relations. It spends longer being planned than executing, and its chosen plan differs between runs. How would you approach this?

level: principalimportance: nice to knowfreq 24%

basics

~20 s

Both symptoms say the query is past the optimizer's exhaustive-search threshold and is being planned by a randomized heuristic. Reduce the effective relation count by decomposing the statement or materializing stable sub-results, make planning happen once via plan reuse or pinning, and fix the statistics that guide whichever search runs.

open as a page