SQL says what you want, not how to get it. Between the query text and the first row returned sits the query planner, also called the optimizer, which chooses how to read each table, in what order to join them, with which join algorithm, and where to sort or aggregate. For a five-table join there are thousands of valid plans, and the best and worst can differ in run time by several orders of magnitude.

This article explains how a cost-based planner makes that choice, from first principles, using PostgreSQL as the concrete example because its internals are documented and its knobs are visible. The same architecture, statistics feeding a cost model that ranks candidate plans produced by a search, underlies Oracle, SQL Server, MySQL and the big-data engines; the Hive cost-based optimizer shows the Calcite variant. By the end you should be able to read a plan, find the estimate that went wrong, and fix the cause rather than the symptom.

Advertisement

The pipeline from SQL to plan

A query passes through a fixed sequence of stages. The parser turns text into a syntax tree. The analyzer resolves names against the catalog, assigns types and checks permissions. The rewriter expands views and applies rules. Only then does the planner see a query tree, which it first simplifies with logical rewrites and then turns into a physical plan.

From SQL text to an executing plan: the planner pipeline and its inputsSQL textclient queryParser + analyzernames, types, treeRewriterviews, rulesLogical rewritespushdown, flattenAccess pathsseq, index, bitmapJoin searchDP or geneticCost modelpages, tuples, opsStatisticspg_stats, extendedCheapest plan treephysical operatorsExecutorpull-based iteratorsPlan cacheprepared statementspathscost?selectivityreuseEstimates flow up from statistics; actual row counts only appear in EXPLAIN ANALYZE and are not fed back automatically.
The planner generates access paths and join orders, asks the cost model to price each candidate using statistics, and hands the cheapest tree to the executor. Prepared statements can reuse a cached plan.

Logical rewrites are transformations that are always, or almost always, beneficial and need no costing: pushing filters down toward the scans, flattening simple subqueries and views into the outer query so their tables can be reordered, turning some IN and EXISTS subqueries into semi-joins, and removing joins whose results are provably unused. Since PostgreSQL 12, a non-recursive common table expression referenced once is inlined unless you write MATERIALIZED; before that, every CTE was an optimisation fence.

Physical planning is where costing starts. For each base table the planner lists access paths: a sequential scan, an index scan per usable index, a bitmap scan that collects matching row locations from one or more indexes and reads heap pages in physical order, and an index-only scan when the index covers every needed column. How indexes attach to storage is covered in database indexing architecture.

The cost model

A plan's cost is an estimate in arbitrary units, anchored so that reading one page sequentially costs 1.0. PostgreSQL's default constants are seq_page_cost = 1.0, random_page_cost = 4.0, cpu_tuple_cost = 0.01, cpu_index_tuple_cost = 0.005 and cpu_operator_cost = 0.0025. Each operator has a formula combining these with estimated page and row counts, and a parent's cost includes its children's.

Take an invented orders table with 1,000,000 rows in 12,500 pages, 80 rows per page, and a query with one filter, WHERE customer_id = 42. A sequential scan reads every page and evaluates the filter on every row:

seq scan cost = pages * seq_page_cost + rows * cpu_tuple_cost + rows * cpu_operator_cost
              = 12,500 * 1.0          + 1,000,000 * 0.01     + 1,000,000 * 0.0025
              = 12,500 + 10,000 + 2,500 = 25,000

If statistics say 1,000 rows match, an index scan fetches those rows through the index. In a simplified model where matching rows are scattered, each is a random heap page read: about 1,000 times 4.0 = 4,000, plus small CPU terms. The index wins easily. The real formula is more subtle: it uses the column's physical correlation, which, if rows for a customer are clustered, makes fetches nearly sequential, and effective_cache_size, which models pages likely to be cached. But the shape of the trade-off is visible even in the simple model: as the number of matching rows rises into the several thousands, random reads approach the cost of reading the whole table in order, and the planner switches first to a bitmap scan, which sorts row locations by page, and then to a sequential scan.

Two consequences matter operationally. The constants describe hardware, so on SSD storage with a mostly cached working set, a random_page_cost closer to 1.1 to 2.0 is common and shifts the crossover toward indexes. And every cost is proportional to estimated row counts, so a wrong row estimate produces a wrong cost however good the formulas are.

Advertisement

Statistics and selectivity estimation

ANALYZE samples each table and stores per-column statistics, visible in the pg_stats view: the fraction of nulls, the estimated number of distinct values, the most common values and their frequencies, a histogram of the remaining values with roughly equal numbers of rows per bucket, and the correlation between value order and physical order. default_statistics_target is 100, which controls the number of most-common values and histogram buckets kept and the sample size; it can be raised per column.

Selectivity, the fraction of rows a predicate keeps, is estimated from these. An equality on a common value uses its stored frequency. An equality on another value spreads the remaining frequency over the remaining distinct values. A range uses the histogram. Then comes the assumption that causes most bad plans: for a = 1 AND b = 2, the planner multiplies the two selectivities as if the columns were independent.

Join sizes are estimated similarly, from the distinct counts and most common values of the join columns on each side. Errors compound: an estimate that is ten times too low at a scan becomes a hundred times too low after two joins, and join algorithms are chosen on those numbers. A nested loop is excellent for 2 outer rows and disastrous for 200,000, which is why hash join exists as the planner's choice for large inputs.

Join search: dynamic programming and its limits

The classical algorithm comes from IBM's System R optimizer, described by Selinger and colleagues in 1979, and PostgreSQL's standard join search follows the same idea. Build the cheapest plan for every single table, then for every pair, then every triple, each time combining the best plans of smaller subsets. Keep more than one plan per subset when a more expensive plan produces rows in a useful order, an 'interesting order', that could avoid a sort later.

# Selinger-style dynamic programming over join subsets (simplified).
# best[S] = cheapest plan producing the join of relation set S,
# kept per interesting order (e.g. sorted on a join or ORDER BY key).
def plan_joins(relations, cost, access_paths, join_methods):
    best = {}
    for r in relations:                               # level 1: base scans
        best[frozenset([r])] = min(access_paths(r), key=cost)
    for size in range(2, len(relations) + 1):         # level k from levels below
        for S in subsets_of_size(relations, size):
            candidates = []
            for left in proper_nonempty_subsets(S):
                right = S - left
                if left not in best or right not in best:
                    continue
                if not connected_by_predicate(left, right):
                    continue                          # avoid cross joins
                for method in join_methods:           # nested loop, hash, merge
                    candidates.append(method(best[left], best[right]))
            if candidates:
                best[frozenset(S)] = min(candidates, key=cost)
    return best[frozenset(relations)]

The number of subsets grows exponentially with the number of tables, so exhaustive search stops being affordable somewhere past a dozen relations. PostgreSQL has two guards. join_collapse_limit and from_collapse_limit, both 8 by default, limit how many relations are flattened into one search problem; beyond that, explicit JOIN syntax order is partly kept. And at geqo_threshold, 12 by default, the planner switches to the genetic query optimizer, a randomised search that finds a good plan rather than the optimal one.

Other systems make different choices. SQL Server and Apache Calcite use Cascades or Volcano-style frameworks, which explore a shared memo of equivalent expressions using transformation rules and prune with cost bounds, and several engines add adaptive execution that re-plans between stages using observed sizes, as Spark EXPLAIN plans show for Spark SQL.

Worked example: a misestimate and its fix

An orders table has 1,000,000 rows with city and zip columns. 'Oakland' has frequency 0.005 and zip 94612 has frequency 0.0004. Zip determines city, so every order in 94612 is also in Oakland: the true result of filtering on both is 0.0004 times 1,000,000, or 400 rows. The planner, assuming independence, multiplies: 0.005 times 0.0004 times 1,000,000 gives 2 rows.

With an estimate of 2, joining to customers by nested loop with an index probe per row looks nearly free. At 400 rows it is still tolerable; on a larger table with the same 200-fold error, the plan does hundreds of thousands of random probes. The symptom in EXPLAIN ANALYZE is the gap between estimated and actual rows on the lowest node:

-- Illustrative, trimmed output before extended statistics
Nested Loop  (cost=0.85..29.38 rows=2 width=96)
             (actual time=0.05..182.4 rows=400 loops=1)
  ->  Index Scan using orders_city_zip_idx on orders o
        (cost=0.42..12.46 rows=2 width=64) (actual rows=400 loops=1)
        Index Cond: ((city = 'Oakland') AND (zip = '94612'))
  ->  Index Scan using customers_pkey on customers c
        (cost=0.43..8.45 rows=1 width=32) (actual rows=1 loops=400)
        Index Cond: (id = o.customer_id)

The fix is to tell the planner about the dependency, not to force a join method. Extended statistics, available since PostgreSQL 10, record functional dependencies, multi-column distinct counts and, since PostgreSQL 12, multi-column most-common-value lists:

-- What the planner knows about a column
SELECT attname, null_frac, n_distinct, most_common_vals, most_common_freqs,
       correlation
FROM pg_stats
WHERE tablename = 'orders' AND attname IN ('city', 'zip');

-- The estimate before and after teaching it the dependency
EXPLAIN SELECT * FROM orders WHERE city = 'Oakland' AND zip = '94612';

CREATE STATISTICS orders_city_zip (dependencies) ON city, zip FROM orders;
ANALYZE orders;

EXPLAIN SELECT * FROM orders WHERE city = 'Oakland' AND zip = '94612';

After ANALYZE, the dependency statistic tells the planner that knowing the zip nearly determines the city, and the combined estimate moves close to 400. The plan is now chosen on the right number, and it stays right as data grows. Read plans bottom up and compare estimated with actual rows at every node; the lowest node with a large ratio is almost always where to look.

Plan caching and parameters

Planning costs CPU, so drivers and applications reuse plans through prepared statements. PostgreSQL plans the first five executions of a prepared statement with the actual parameter values, called custom plans, and then considers a generic plan that ignores the values, switching to it if its estimated cost is not meaningfully worse than the average custom plan. Since PostgreSQL 12, plan_cache_mode can force one behaviour or the other.

Generic plans go wrong on skewed data. If one tenant owns half the rows and most own a few, a generic plan chosen for the typical tenant uses an index scan that is ruinous for the large one. Symptoms are a statement that is fast in manual testing, where values are literals, and slow from the application. Setting plan_cache_mode = force_custom_plan for that role or session, or splitting the query for the known heavy values, fixes it at the cost of planning time per execution.

Failure modes

  • Stale statistics. A bulk load or delete runs before autovacuum's analyze threshold is crossed; estimates describe yesterday's table. Run ANALYZE explicitly after large loads.
  • Correlated predicates. The independence assumption underestimates. Use extended statistics on the column groups your filters actually combine.
  • Opaque expressions. Functions of columns, such as lower(email) = 'x', have no statistics and get default selectivities. An expression index also collects statistics for the expression.
  • Type mismatches. Comparing a text column with a numeric parameter or across collations can make an index unusable. Match types in the query.
  • LIMIT with ORDER BY. The planner may pick an index in sort order, expecting to find matches early; if matching rows are rare or at the end, it walks most of the index.
  • Too many joins. Past the collapse limits or the genetic optimizer threshold, plans depend on join syntax order or random search. Simplify the query or raise the limits for that session.

Operating the planner

Use EXPLAIN (ANALYZE, BUFFERS) on a representative parameter set, and remember that it executes the statement, so wrap data-modifying statements in a transaction you roll back. Enable pg_stat_statements to find the statements that consume the most total time, and log slow plans automatically with the auto_explain module. Planner switches such as enable_nestloop = off are diagnostic tools for proving that another plan would be faster; they are not production fixes, because they affect every query in the session. Index design interacts with all of this, and PostgreSQL index types covers which structures each predicate can use.

The trade-off at the heart of every planner is planning time against plan quality. Exhaustive search, rich statistics and per-execution planning give better plans and cost CPU and memory on every query; caches, limits and heuristics are cheap and occasionally very wrong. Your job is mostly to keep the inputs accurate so that the cheap path stays right.

What to do next

  1. Pick your ten most expensive statements from pg_stat_statements and run EXPLAIN (ANALYZE, BUFFERS) on each.
  2. At every node, compare estimated and actual rows; note the lowest node where they differ by ten times or more.
  3. For misestimates on combined filters, create extended statistics on those columns and re-analyze.
  4. Check that bulk loads are followed by ANALYZE, and raise statistics targets on skewed columns.
  5. Review random_page_cost against your storage and cache hit ratio, and test changes per session first.
  6. For parameterised queries on skewed data, compare custom and generic plans and set plan_cache_mode where needed.
Key takeaway: A query planner turns SQL into a plan by generating candidate access paths and join orders, pricing each with a cost model, and keeping the cheapest. The cost formulas are simple; the hard part is estimating row counts, which depend on statistics and on assumptions such as column independence. Most bad plans are bad estimates, visible as a gap between estimated and actual rows in EXPLAIN ANALYZE. Fix the inputs with fresh and extended statistics, realistic cost constants and sensible plan caching, and use planner switches only to diagnose.