I think I should preface this one with an apology, it’s a bit long. Other than that I hope you enjoy it.

Preface

Around spring of last year I spent some of my time following the CMU OPTIMIZE! course on query optimization. Not all the lectures ended up being filmed because of Andy’s illness but the ones that were are possibly the best I’ve seen, not only because there isn’t really another course on the subject anywhere except TUM’s, but also because it felt like someone turned on the wipers while it was raining. A lot of the blurriness I had on the subject, mostly from going through TUM’s lectures and papers, was gone and I finally had a good mental model of it.

I was initially going to write a sort of literature review of query optimization at large and talk a bit about things like Volcano and Cascades. But I quickly realized that it would take a few hundred pages, someone is already doing that, and I highly encourage you to check out the full course.

Instead I decided to focus on join ordering since I feel it’s the most relevant part of query optimization that doesn’t get the entry level treatment. Specifically I want to focus on the German lineage of dynamic programming approaches. Don’t worry, I also think cardinality estimation isn’t that cute and I plan to cover it sometime later.

This started as a notes document while I was working through the course, initially I thought I’d end up with a couple of pages I could refer back to but it quickly grew to about 46 pages. I realized it would be better shared here, and at the same time it could use a lot of editing.

The series is split into 6 parts that build on each other. With enough background some parts can be read on their own but I would recommend following them in order. I was initially planning to release everything as one long post but splitting it up seemed like a better idea for everyone’s attention span, and I will probably publish the whole thing in one piece once the series is complete.

We will focus on the join ordering line of research from the Munich school of query optimization and limit ourselves to 4 papers that I will cover in as much depth as I can, while referring back to other papers and approaches along the way.

The papers are:

-

Moerkotte & Neumann, “Analysis of Two Existing and One New Dynamic Programming Algorithm for the Generation of Optimal Bushy Join Trees without Cross Products” (VLDB 2006) - which analyzes DPsize and DPsub, defines the csg-cmp-pair as the fundamental unit of work and introduces DPccp, the first DP algorithm that touches each unit exactly once.

-

Moerkotte & Neumann, “Dynamic Programming Strikes Back” (SIGMOD 2008) - which generalizes DPccp from graphs to hypergraphs, yielding DPhyp, the algorithm that handles complex predicates and crucially the reordering restrictions imposed by outer joins.

-

Neumann & Radke, “Adaptive Optimization of Very Large Join Queries” (SIGMOD 2018) - which asks what to do when exact DP is hopeless and you have hundreds or thousands of relations and answers with a strategy built from IKKBZ, linearized DP, and GOO.

-

Birler & Neumann, “Efficient Enumeration of the Complete Join Search Space” (DBPL 2025) - which closes a fifteen-year-old gap: hypergraph-based conflict detection for non-inner joins was efficient but incomplete (CD-A), or complete but slow (CD-C). CD-E is both, by pruning reordering restrictions that a cross-product-avoiding enumerator could never violate anyway.

So without making this post any longer than it already is, I leave you with Part 1.

Join Ordering, Affairs.

In the hidden lore of SQL books (and the SQL standard, which I’ll admit I have never fully read) nothing is said about the order in which joins happen, only about what the result has to look like. How to get there is left to the query engine, more specifically to the planner or optimizer, and its job is to pick a join tree: a binary tree whose leaves are the base relations and whose inner nodes are join operators. All trees end up computing the same result for inner joins, but their costs can differ wildly because the sizes of the intermediate results differ wildly[1].

Let’s consider a starting example, we have three relations with |R₀| = |R₁| = 1,000,000 and |R₂| = 10, a highly

selective predicate between R₀ and R₂, another between R₂ and R₁, and no predicate between R₀ and R₁. The plan

(R₀ ⋈ R₂) ⋈ R₁ first shrinks R₀ down to a handful of rows and finishes cheaply. The plan (R₀ ⋈ R₁) ⋈ R₂ must first

build a cross-product-like monster of up to 10¹² tuples before filtering out on the R₂ predicate, same result but ten

orders of magnitude apart[2].

If you want some data on how much this matters in real systems have a look at “Query optimization through the looking glass”, where the planner’s ordering choice alone swings execution times by 100x and more.

One might think, well can’t we just enumerate all the orderings and cost each one ?. In the example above the selectivity of the predicates already tells us which plan is cheaper, and for three relations that works fine. The problem is how many orderings there are to look at[3]. We know that the number of binary trees over n leaves is the Catalan number C(n−1), and since the leaves can be permuted the number of distinct join trees is

from math import factorial

def distinct_join_trees(n: int) -> int:

"""(2(n-1))! / (n-1)! -- the number of distinct join trees over n relations."""

return factorial(2 * (n - 1)) // factorial(n - 1)

for n in (4, 5, 15):

print(n, distinct_join_trees(n))which is 120 for n = 4, 1,680 for n = 5 and about 3.5 · 10¹⁸ for n = 15. Join ordering with a real cost function is NP-hard in general, even if you restrict yourself to left-deep trees with cross products. But a lot of NP-hard problems sit within the solvable fast enough search space and in the case of join ordering most modern systems can solve moderate instances exactly via dynamic programming. In Part 4 when we will look at IKKBZ we will see how in some cases you can relax the hardness assumptions and get optimal left-deep orderings in polynomial time.

The idea behind the dynamic programming is simple enough to state now. Take the best plan for all the relations in a query, its top join splits them into two sets and each side is itself planned in some way. If there was a cheaper plan for one of the sides you could swap it in and get a cheaper plan overall, so the best plan for a set is always built from the best plans of two smaller sets. Instead of looking at whole trees we look at sets of relations, starting from single relations and working our way up, and for every set we remember the cheapest plan we found. The rest of this part is about what it takes to make that work, which sets are worth planning at all, how to represent them, how to estimate their size and cost, and how much work the DP can’t avoid doing.

Left-deep, right-deep, zig-zag and bushy trees

Join trees come in a few shapes. A tree is left-deep when the right child of every join is a base relation, so the tree is a spine hanging to the left and the whole plan is described by a permutation of the relations. Right-deep is the mirror image. Zig-zag trees allow the base relation to sit on either side of each join, which matters for physical operators since in a hash join you want the smaller input on the build side. A tree with no restriction on its shape is bushy, both inputs of a join can themselves be the result of other joins.

left-deep: ⋈ right-deep: ⋈

/ \ / \

⋈ R₃ R₀ ⋈

/ \ / \

⋈ R₂ R₁ ⋈

/ \ / \

R₀ R₁ R₂ R₃

zig-zag: ⋈ bushy: ⋈

/ \ / \

⋈ R₃ ⋈ ⋈

/ \ / \ / \

R₂ ⋈ R₀ R₁ R₂ R₃

/ \

R₀ R₁If we count them there are n! left-deep trees against n! · C(n−1) bushy ones, so the bushy space is larger by a factor of

which grows exponentially with n.

Why bushy trees can be strictly better

System R only ever searched left-deep trees and for a long time the common wisdom was that this is good enough. It

isn’t, and there are two reasons for it. The first one is the size of intermediate results. Consider a chain query

R₀ — R₁ — R₂ — R₃ where the predicates on (R₀,R₁) and (R₂,R₃) are very selective but the one in the middle on (R₁,R₂) is

not. The bushy plan (R₀ ⋈ R₁) ⋈ (R₂ ⋈ R₃) computes two tiny results and joins them once through the bad predicate,

while every left-deep plan has to drag at least one large intermediate result through the middle, e.g.

((R₀ ⋈ R₁) ⋈ R₂) ⋈ R₃ materializes the big R₀R₁R₂ result. On shapes like this the best bushy plan can beat the best

left-deep plan by an arbitrarily large factor, and snowflake schemas where each dimension has its own chain of selective

sub-dimensions are where you run into this in practice. The second reason is parallelism, the two subtrees of a bushy

join are independent and can run concurrently whereas a left-deep spine is a chain of dependencies.

The price for this is a much bigger search space, which is why the enumeration algorithms we will look at are engineered so carefully. Planning sets instead of trees is what keeps it manageable, a DP restricted to left-deep trees only ever splits one relation off a set while a bushy one considers every way of splitting it in two, but both fill the same table.

The query graph

To reason about which trees are worth building we need the query graph G = (V, E). It has one node per base relation

and one edge per binary join predicate, so R.a = S.b becomes the edge {R, S} annotated with the estimated

selectivity of the predicate. There are four shapes that keep coming back in the literature because they cover

the range from easy to hard:

chain: R₀ — R₁ — R₂ — R₃ star: R₁

|

cycle: R₀ — R₁ R₃ — R₀ — R₂

| | |

R₃ — R₂ R₄ (center joins every satellite)

clique: every pair connected (n(n−1)/2 edges)Chains show up with foreign-key pipelines, stars with fact and dimension tables, and cliques when many relations join on one shared attribute and the optimizer adds the transitive closure of the equalities. Real queries usually sit somewhere between a chain and a star, e.g. snowflakes or trees with a few cycles.

The reason we care about the graph is one observation, a join between two subplans whose relation sets have no edge between them is a cross product. Cross products are almost never a good idea since they multiply cardinalities with a selectivity of 1, so most optimizers refuse to introduce them unless the query graph itself is disconnected. Once you rule out cross products join enumeration becomes a graph problem, the only sets of relations worth planning are the ones that form connected subgraphs, and those are the only sets that ever make it into the DP table.

Our three relations from the beginning form a chain R₀ — R₂ — R₁, there is no edge between R₀ and R₁ so the bad plan

(R₀ ⋈ R₁) ⋈ R₂ starts with exactly such a cross product. {R₀, R₁} isn’t a connected subgraph, so an optimizer that

only plans connected subgraphs never builds that plan in the first place.

Sets of relations as bitsets

We now have the shape of a recipe, we want to do some form of dynamic programming over sets of relations. If you’ve done

competitive programming style problems before then what will follow will be familiar.

We want to think about sets of relations in terms of subsets and we want our solving logic to be fast, our goal is to

build a table of sets and look them up as we split and merge them. An efficient way to do this is to use bitmaps or bitsets

to represent sets. The set of all relations {R₀, R₁, R₂} will be represented by [1, 1, 1] the subset {R₀, R₂} will be

[1, 0, 1] and {R₁, R₂} will be [0, 1, 1] and so on and so on. This approach is very popular when doing dynamic

programming problems over sets and it allows you to write fairly compact union, find, merge, lookup code since if you

fix your working set n all your code will be bitwise operations or combination of bitwise operations. Our solver will

follow this same pattern.

/// A set of relations, represented as a 64-bit bitset.

#[derive(Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Debug)]

pub struct NodeSet(pub u64);

impl NodeSet {

pub const EMPTY: NodeSet = NodeSet(0);

pub fn single(i: usize) -> NodeSet { NodeSet(1u64 << i) }

/// {v_0, ..., v_i} — the "B_i" prefix set from the DPccp paper.

pub fn prefix(i: usize) -> NodeSet { NodeSet((1u64 << (i + 1)) - 1) }

pub fn full(n: usize) -> NodeSet {

NodeSet(if n == 64 { !0 } else { (1u64 << n) - 1 })

}

pub fn is_empty(self) -> bool { self.0 == 0 }

pub fn contains(self, i: usize) -> bool { self.0 & (1u64 << i) != 0 }

pub fn is_subset_of(self, other: NodeSet) -> bool { self.0 & !other.0 == 0 }

pub fn intersects(self, other: NodeSet) -> bool { self.0 & other.0 != 0 }

pub fn union(self, other: NodeSet) -> NodeSet { NodeSet(self.0 | other.0) }

pub fn intersect(self, other: NodeSet) -> NodeSet { NodeSet(self.0 & other.0) }

pub fn minus(self, other: NodeSet) -> NodeSet { NodeSet(self.0 & !other.0) }

pub fn len(self) -> u32 { self.0.count_ones() }

/// Index of the minimal element.

pub fn min_elem(self) -> usize {

debug_assert!(!self.is_empty());

self.0.trailing_zeros() as usize

}

/// Iterate over the elements (ascending).

pub fn iter(self) -> impl Iterator<Item = usize> {

let mut bits = self.0;

std::iter::from_fn(move || {

if bits == 0 { None } else {

let i = bits.trailing_zeros() as usize;

bits &= bits - 1; // clear lowest set bit

Some(i)

}

})

}

/// Iterate over all NON-EMPTY subsets of `self`, in ASCENDING numeric

/// order. Since A ⊂ B implies A < B numerically, this guarantees that

/// every subset is produced before any of its supersets — a property the

/// DP algorithms rely on (subplans must exist before they are combined).

/// Trick: sub = (sub - mask) & mask walks all submasks ascending.

pub fn nonempty_subsets(self) -> impl Iterator<Item = NodeSet> {

let mask = self.0;

let mut sub = 0u64;

let mut done = mask == 0;

std::iter::from_fn(move || {

if done { return None; }

sub = sub.wrapping_sub(mask) & mask;

if sub == mask { done = true; }

Some(NodeSet(sub))

})

}

}Submask enumeration

One particular benefit of this approach is that it allows us to have very fast subset enumeration in ascending or

order via masks. sub = (sub - mask) & mask visits every submask of mask in increasing numeric order at

a constant cost per submask, 2^k iterations for a k-bit mask.

The order here matters, if A is a strict subset of B then A < B numerically, so when your dynamic programming loop walks

submasks in increasing order every strict subset of the current set has already been planned. Get this wrong, going in

decreasing order or in whatever order a hash map gives you, and the DP will look up table entries that don’t exist yet.

I know this because I got it wrong during development and the three algorithms from the next part immediately started disagreeing with each other.

The query graph in Rust

/// A (simple) query graph: relations are nodes, binary join predicates are

/// edges annotated with a selectivity.

#[derive(Clone, Debug)]

pub struct QueryGraph {

pub n: usize,

/// cardinalities |R_i|

pub card: Vec<f64>,

/// adjacency[i] = set of neighbors of node i

pub adjacency: Vec<NodeSet>,

/// edges (u, v, selectivity) with u < v

pub edges: Vec<(usize, usize, f64)>,

}

impl QueryGraph {

pub fn new(card: Vec<f64>) -> QueryGraph {

let n = card.len();

QueryGraph { n, card, adjacency: vec![NodeSet::EMPTY; n], edges: vec![] }

}

pub fn add_edge(&mut self, u: usize, v: usize, sel: f64) {

let (u, v) = if u < v { (u, v) } else { (v, u) };

self.adjacency[u] = self.adjacency[u].union(NodeSet::single(v));

self.adjacency[v] = self.adjacency[v].union(NodeSet::single(u));

self.edges.push((u, v, sel));

}

/// Neighborhood of a set S: all nodes adjacent to S, excluding S itself

/// and an exclusion set X.

pub fn neighbors(&self, s: NodeSet, x: NodeSet) -> NodeSet {

let mut acc = NodeSet::EMPTY;

for v in s.iter() {

acc = acc.union(self.adjacency[v]);

}

acc.minus(s).minus(x)

}

/// Is the subgraph induced by S connected? (BFS over bitsets)

pub fn is_connected(&self, s: NodeSet) -> bool {

if s.is_empty() { return false; }

let mut reach = NodeSet::single(s.min_elem());

loop {

let grown = reach.union(self.neighbors(reach, NodeSet::EMPTY).intersect(s));

if grown == reach { break; }

reach = grown;

}

reach == s

}

/// Is there at least one edge with one endpoint in S1 and the other in S2?

pub fn connected_pair(&self, s1: NodeSet, s2: NodeSet) -> bool {

!self.neighbors(s1, NodeSet::EMPTY).intersect(s2).is_empty()

}

/// Estimated cardinality of the join of all relations in S:

/// |⋈ S| = ∏ card(R_i) · ∏ sel(e) over edges internal to S.

pub fn cardinality(&self, s: NodeSet) -> f64 {

let mut c = 1.0;

for v in s.iter() { c *= self.card[v]; }

for &(u, v, sel) in &self.edges {

if s.contains(u) && s.contains(v) { c *= sel; }

}

c

}

}The cardinality function is the classical independence model from Selinger-style estimation, the size of a join

result is the product of the base cardinalities and the selectivities of every predicate applied. Every DP we

will look at relies on one property of it, the estimate for a set S depends only on S and not on the join tree used to produce it.

All plans for the same S produce the same number of rows, so “the cost of the best plan for S” is well defined and

Bellman’s principle of optimality applies. Real optimizers replace this model with histograms, sketches and

sampling[4], the estimates get better but as long as the cardinality stays a function of the set the enumeration

machinery doesn’t change.

Here is the example from the beginning as a query graph. I picked the selectivities so that every row of R₂ matches exactly one row of R₀ and one row of R₁, which gives 10⁻⁶ for both predicates:

/// Builds a set from a list of relation indices.

fn set(rels: &[usize]) -> NodeSet {

rels.iter().fold(NodeSet::EMPTY, |acc, &i| acc.union(NodeSet::single(i)))

}

/// The three relations from the beginning: |R₀| = |R₁| = 1,000,000 and

/// |R₂| = 10, joined as the chain R₀ — R₂ — R₁.

fn example() -> QueryGraph {

let mut g = QueryGraph::new(vec![1e6, 1e6, 10.0]);

g.add_edge(0, 2, 1e-6);

g.add_edge(2, 1, 1e-6);

g

}

#[test]

fn example_cardinalities() {

let g = example();

// no predicate applies between R₀ and R₁, it's a cross product

assert!(!g.connected_pair(set(&[0]), set(&[1])));

assert_eq!(g.cardinality(set(&[0, 1])), 1e12);

// the selective joins keep everything down to R₂'s 10 rows

assert_eq!(g.cardinality(set(&[0, 2])).round(), 10.0);

assert_eq!(g.cardinality(set(&[0, 1, 2])).round(), 10.0);

}The estimate for {R₀, R₁, R₂} comes out as 10 rows no matter which of the two plans produces it, which is the property the DP needs. What differs between the plans is what they build on the way there, and that’s what the cost model is for.

Plans and Cost Models

A plan can either be a base relation or a join of two plans and the DP table maps each set of relations to the cheapest plan found for it so far:

use std::collections::HashMap;

#[derive(Clone, Debug)]

pub enum Plan {

Leaf(usize),

Join(Box<Plan>, Box<Plan>),

}

#[derive(Clone, Debug)]

pub struct Entry {

pub set: NodeSet,

pub card: f64,

pub cost: f64,

pub plan: Plan,

}

pub type DpTable = HashMap<NodeSet, Entry>;

/// The cost of a join is the sum of the costs of its inputs plus the

/// cardinality of its output. Leaves cost 0.

pub fn combine(graph: &QueryGraph, left: &Entry, right: &Entry) -> Entry {

let set = left.set.union(right.set);

let card = graph.cardinality(set);

Entry {

set,

card,

cost: left.cost + right.cost + card,

plan: Plan::Join(Box::new(left.plan.clone()), Box::new(right.plan.clone())),

}

}

fn seed_leaves(graph: &QueryGraph) -> DpTable {

let mut t = DpTable::new();

for i in 0..graph.n {

let s = NodeSet::single(i);

t.insert(s, Entry { set: s, card: graph.card[i], cost: 0.0, plan: Plan::Leaf(i) });

}

t

}

fn consider(table: &mut DpTable, cand: Entry) {

match table.get(&cand.set) {

Some(old) if old.cost <= cand.cost => {}

_ => { table.insert(cand.set, cand); }

}

}The cost of a plan Cout is the sum of the sizes of all its intermediate results. This is the standard cost model in most papers

I’ve found. It’s symmetric since swapping the inputs of a join doesn’t change the cost, it’s monotone, and it has the

adjacent sequence interchange (ASI) property on left-deep plans which we will need later for IKKBZ. Production cost

models do a lot more, they distinguish the build side of a hash join from the probe side, account for memory and so on

but the enumeration algorithms don’t care. They build a candidate with combine and hand it to consider, which only

keeps it if it’s cheaper than what the table already has for that set. With Cout one candidate per split is enough

since the cost is symmetric, with an asymmetric cost model you would call combine for both orders of the inputs and

let consider keep the cheaper one.

Assigning a cost to the two plans from the beginning, (R₀ ⋈ R₂) ⋈ R₁ builds two intermediate results of 10 rows each and

costs 20, while (R₀ ⋈ R₁) ⋈ R₂ pays for the 10¹² rows of the cross product before getting down to the same 10:

#[test]

fn example_plan_costs() {

let g = example();

let t = seed_leaves(&g);

let leaf = |i: usize| &t[&NodeSet::single(i)];

let good = combine(&g, &combine(&g, leaf(0), leaf(2)), leaf(1));

let bad = combine(&g, &combine(&g, leaf(0), leaf(1)), leaf(2));

assert_eq!(good.cost.round(), 20.0);

assert_eq!(bad.cost.round(), 1e12 + 10.0);

}The csg-cmp-pair

The most important idea in the VLDB 2006 paper is giving a name to the work that can’t be avoided. A pair of relation sets (S₁, S₂) is a csg-cmp-pair (connected subgraph, connected complement) when

- S₁ induces a connected subgraph of the query graph,

- S₂ induces a connected subgraph, disjoint from S₁, and

- at least one edge connects S₁ and S₂.

Every inner node of a join tree without cross products is such a pair, the node joins plan(S₁) with plan(S₂), and the

other way around every csg-cmp-pair is the top-level split of some valid tree. So any bottom-up optimizer that finds the

best bushy plan without cross products has to look at every csg-cmp-pair at least once, which makes their number #ccp

a lower bound on the work of enumeration.

Our example has four of them: ({R₀}, {R₂}), ({R₂}, {R₁}), ({R₀}, {R₂, R₁}) and ({R₀, R₂}, {R₁}). ({R₀}, {R₁}) isn’t one because no edge connects the two, and neither is ({R₂}, {R₀, R₁}) because {R₀, R₁} isn’t connected. The plan that cost us 10¹² would need one of those two splits, so a DP that only looks at csg-cmp-pairs never gets to consider it.

For larger queries we need to count, and the paper gives closed forms for the four shapes from earlier[5]:

We can compute the counts of the csg-cmp-pairs by brute force using the code we have so far to check the closed form formulas.

/// Counts the unordered csg-cmp-pairs of `g` by brute force: every pair of

/// disjoint, connected relation sets with at least one edge between them.

fn count_ccps(g: &QueryGraph) -> u64 {

let all = NodeSet::full(g.n);

let mut count = 0;

for s1 in all.nonempty_subsets() {

if !g.is_connected(s1) { continue; }

for s2 in all.minus(s1).nonempty_subsets() {

// s1 < s2 so each unordered pair is counted once

if s1 < s2 && g.is_connected(s2) && g.connected_pair(s1, s2) {

count += 1;

}

}

}

count

}

fn chain(n: usize) -> QueryGraph {

let mut g = QueryGraph::new(vec![100.0; n]);

for i in 0..n - 1 { g.add_edge(i, i + 1, 0.1); }

g

}

fn cycle(n: usize) -> QueryGraph {

let mut g = chain(n);

g.add_edge(n - 1, 0, 0.1);

g

}

fn star(n: usize) -> QueryGraph {

let mut g = QueryGraph::new(vec![100.0; n]);

for i in 1..n { g.add_edge(0, i, 0.1); }

g

}

fn clique(n: usize) -> QueryGraph {

let mut g = QueryGraph::new(vec![100.0; n]);

for i in 0..n {

for j in i + 1..n { g.add_edge(i, j, 0.1); }

}

g

}

#[test]

fn ccp_counts_match_closed_forms() {

for n in 3..=10u64 {

let k = n as usize;

assert_eq!(count_ccps(&chain(k)), (n * n * n - n) / 6);

assert_eq!(count_ccps(&cycle(k)), (n * n * n - 2 * n * n + n) / 2);

assert_eq!(count_ccps(&star(k)), (n - 1) * (1 << (n - 2)));

assert_eq!(count_ccps(&clique(k)), (3u64.pow(n as u32) - (1 << (n + 1)) + 1) / 2);

}

}

#[test]

fn example_ccps() {

assert_eq!(count_ccps(&example()), 4);

}This holds for all four shapes up to 10 relations. Do note that the brute force is expensive, it looks at every subset and then at every subset of its complement, which is 3ⁿ candidate pairs for n relations, while a chain only has (n³ − n) / 6 real ones. Almost all of the work goes into pairs that fail one of the three conditions.

What the table above tells us is that the shape of the query graph matters more than the number of relations. With 20 relations a chain has 1,330 pairs which is nothing, a star has 19 · 2¹⁸ ≈ 5 · 10⁶ which is fine, and a clique has about 1.7 · 10⁹ which is where things start to hurt. It also tells us what an optimal algorithm looks like, one that spends constant time per csg-cmp-pair and doesn’t enumerate anything else. DPsize and DPsub don’t do that, DPccp does, and that’s where the next part picks up.

References

- Guido Moerkotte, Thomas Neumann, Analysis of Two Existing and One New Dynamic Programming Algorithm for the Generation of Optimal Bushy Join Trees without Cross Products, VLDB 2006.

- Guido Moerkotte, Thomas Neumann, Dynamic Programming Strikes Back, SIGMOD 2008.

- Thomas Neumann, Bernhard Radke, Adaptive Optimization of Very Large Join Queries, SIGMOD 2018.

- Altan Birler, Thomas Neumann, Efficient Enumeration of the Complete Join Search Space, DBPL 2025.

- Viktor Leis et al., How Good Are Query Optimizers, Really?, PVLDB 9.

- Viktor Leis et al., Query Optimization Through the Looking Glass, and What We Found Running the Join Order Benchmark, The VLDB Journal 27.

- Guido Moerkotte, Building Query Compilers.

-

A sufficient mental model to use is that of the HashJoinoperator, usually aHashJoinwill pick the smaller table as the build candidate since the build side ideally resides in memory during the probe phase but that heuristics only considers the two input relations at the physical level, if the logical level doesn’t optimize for ordering then the largest table built from the first ordering will end up being the “build side” and you will pay the cost in cycles of probing it even if the right side is smaller. ↩

-

Magnitude here is in terms of rows (tuples) in the intermediate results, how many bytes that ends up being in memory, or on disk in case of spillage, depends on how wide the rows are. ↩

-

I am omitting the fact that cardinality estimation simplifies the ordering problem when cardinalities are accurate but cardinalities are rarely accurate. ↩

-

Estimation error is empirically a bigger source of bad plans than an incomplete enumeration, that’s the other headline result of the Join Order Benchmark paper. I will cover cardinality estimation separately, my goal for now is to outline the join enumeration side only. ↩

-

These count each unordered pair once, which is also what the counting code below does. The paper counts ordered pairs so its numbers are exactly twice these. ↩