Query Optimization and Execution

Query Optimization and Execution

Definition: The process by which a relational database engine turns a declarative SQL query into an efficient physical execution plan, by parsing it, evaluating alternative strategies against table statistics, and choosing the cheapest one it can find.

How It Works

  • Parse: SQL text is checked for syntax and resolved against the schema (do these tables/columns exist, are the types compatible), producing a logical query tree.
  • Rewrite: the logical tree is transformed by rule-based simplifications, view expansion, predicate pushdown (moving WHERE filters as close to the data source as possible), constant folding.
  • Plan generation: the Cost-Based Optimizer (CBO) enumerates candidate physical plans, different join orders, different access methods per table, and estimates the cost of each using table statistics (row counts, value distributions, index selectivity).
  • Execute: the chosen plan runs as a tree of physical operators, each pulling rows from its children, scans at the leaves, joins and aggregates further up, until the final result streams back to the client.

Under the Hood

A query joining orders and customers, showing the pipeline stages and the shape of the resulting plan tree:

Worked example 1: choosing an access method

  • Given: orders has 50 million rows, an index exists on customer_id, and the query is SELECT * FROM orders WHERE customer_id = 42.
  • Step: the optimizer checks index statistics and estimates this customer has roughly 12 orders out of 50 million, a highly selective filter.
  • Step: it compares the cost of an Index Scan (a handful of page reads) against a Sequential Scan (reading all 50 million rows).
  • Answer: it picks the Index Scan. If the same query filtered WHERE status != 'cancelled' and 95% of rows matched, it would likely pick the Sequential Scan instead, an index scan touching most of the table plus random I/O per row is often slower than one fast sequential sweep.

Worked example 2: join algorithm selection

  • Given: joining a 500-row regions table against a 200-million-row orders table on region_id.
  • Step: the optimizer estimates a Nested Loop Join (for each of the 500 regions, probe an index on orders.region_id) would cost roughly 500 index probes, very cheap since one side is tiny.
  • Step: it compares that against a Hash Join (build an in-memory hash table from the smaller regions side, then stream orders through it once) and a Merge Join (both sides pre-sorted, walked in lockstep).
  • Answer: with one side this small and an index available, Nested Loop wins. If both tables were large and unsorted, Hash Join usually wins; if both were already sorted on the join key (e.g. from an earlier ORDER BY or index), Merge Join usually wins.

Worked example 3: stale statistics cause a bad plan

  • Given: a table’s statistics were last collected when it had 1,000 rows; it now has 50 million, but ANALYZE/statistics refresh hasn’t run.
  • Step: the optimizer still believes the table is small and estimates a Sequential Scan is cheap, or worse, chooses a Nested Loop Join expecting few rows on both sides.
  • Answer: the query runs orders of magnitude slower than it should. Running ANALYZE (PostgreSQL) or UPDATE STATISTICS (SQL Server) refreshes the row count and value distribution histograms, and the same query plan changes to an Index Scan or Hash Join immediately, no query rewrite needed.

Worked example 4: join order matters more than join algorithm

  • Given: a three-way join of orders (200M rows), customers (2M rows filtered down to 500 by WHERE region = 'EU'), and order_items (600M rows).
  • Step: if the optimizer joins orders to order_items first (both huge, no filter applied yet), it builds an enormous intermediate result before ever applying the selective customers filter.
  • Step: if instead it joins the small, filtered customers set to orders first, cutting orders down to a few thousand matching rows, then joins that small intermediate result to order_items, every subsequent step operates on a tiny working set.
  • Answer: the second join order can be orders of magnitude faster despite using the exact same join algorithms, this is why the CBO spends most of its search effort on join ordering, not just algorithm choice, and why hinting a bad join order can hurt more than picking a suboptimal algorithm.

Reading an Execution Plan

Plan trees are read bottom-up: leaves are where data enters (scans), and each parent operator consumes its children’s output.

  • Estimated vs actual rows: EXPLAIN alone shows what the optimizer expected; EXPLAIN ANALYZE shows what actually happened. A large gap between the two (say, estimated 10 rows, actual 2 million) is the single most common signal of a statistics or cardinality-estimation problem.
  • Cost units: the numbers next to each node (e.g. cost=0.42..8.44) are the optimizer’s internal, unitless cost estimate, arbitrary page-reads-and-CPU-cycles currency used only to compare plans against each other, not wall-clock time.
  • Loops: in a Nested Loop Join, watch the “loops” count. If the inner side is scanned 50,000 times because the outer side has 50,000 rows, an inner sequential scan that looked cheap in isolation becomes catastrophic multiplied by the loop count.

Why It Matters

The same SQL query can run in milliseconds or minutes depending entirely on which physical plan the optimizer picks, application code never has to change. Understanding this layer is what separates “the query is slow, let’s add caching” from “the query is slow because the planner picked a Nested Loop Join over 10 million rows, let’s fix the statistics or add an index,” the second fix is usually cheaper and more durable.

  • It turns performance debugging into an evidence-based exercise: EXPLAIN ANALYZE output is a concrete artifact to reason from, instead of guesswork about what “feels slow.”
  • It’s often far cheaper than scaling hardware. A missing index or a stale statistics table can make a query 100x slower than it needs to be, no amount of additional CPU or memory fixes a fundamentally bad plan as reliably as fixing the plan itself.
  • It’s a portable skill across engines: the concepts, cost estimation, join algorithms, access paths, transfer directly between PostgreSQL, MySQL, SQL Server, and Oracle even though the exact syntax for inspecting plans differs.

Common Pitfalls

  • Outdated table statistics causing the optimizer to pick sequential scans over a perfectly good index, or the reverse, treating a low-cardinality index as far more selective than it actually is.
  • Wrapping an indexed column in a function (WHERE YEAR(created_at) = 2024) which usually prevents the optimizer from using an index on created_at at all, unless a matching expression index exists.
  • SELECT * when only two columns are needed, which can prevent an index-only scan and force the engine back to the table’s heap for every row.
  • Assuming EXPLAIN (the plan) tells the whole story. EXPLAIN ANALYZE actually runs the query and reports real row counts and timing per step, which is what you need to spot a bad cardinality estimate.
  • Over-trusting the optimizer on very large, correlated WHERE clauses (e.g. city = 'X' AND zip = 'Y' where the two are not independent), the CBO usually assumes independence between predicates unless extended statistics are defined, and can badly misjudge selectivity as a result.
  • Forcing join order or index choice with hints as the first fix instead of the last resort. Hints freeze a decision the optimizer would otherwise re-evaluate as data grows, and a plan hinted for today’s data distribution can become the wrong plan next year.
  • Ignoring parameter sniffing: a query plan cached for one parameter value (a highly selective WHERE customer_id = 42) can get reused for a very different value (a common customer_id with a million rows), producing wildly inconsistent performance for the “same” query.
  • Not accounting for the planning cost itself. Extremely complex queries with many joins can spend measurable time just searching the plan space, most optimizers cap search effort (e.g. Postgres’s join_collapse_limit), which can itself produce a worse plan.

Optimizer Types

  • Rule-Based Optimizer (RBO): applies a fixed set of heuristics (always prefer an index if one exists, always push filters down) regardless of actual data. Simple and predictable, but blind to data skew.
  • Cost-Based Optimizer (CBO): the standard in modern engines. Uses collected statistics, row counts, histograms, distinct value counts, to estimate the cost of each candidate plan and pick the cheapest. Requires accurate, current statistics to work well.
  • Adaptive/runtime optimizers: some engines (SQL Server’s adaptive joins, Oracle’s adaptive query optimization) can revise a plan mid-execution if early results reveal the initial cardinality estimate was badly wrong, closing part of the estimate-vs-actual gap without a full re-plan.

Comparison

Choosing between join algorithms and scan types is not something application code controls directly, it falls entirely out of the optimizer’s cost estimates, but knowing the trade-offs below is what makes an EXPLAIN ANALYZE output interpretable rather than a wall of unfamiliar node names.

Join algorithmBest whenMemory needsComplexity
Nested LoopOne side is small, or an index exists on the join keyMinimalSimple, O(n × m) worst case
Hash JoinBoth sides large, unsorted, equality joinNeeds to build a hash table for the smaller sideModerate, roughly O(n + m)
Merge JoinBoth sides already sorted on the join keyMinimal if pre-sorted, else needs a sort stepEfficient once sorted, O(n + m)
Scan typeReadsBest when
Sequential ScanEvery row in the tableFilter matches most rows, or table is small
Index ScanIndex pages, then row lookupsFilter is selective, few matching rows
Index-Only ScanIndex pages only, no table lookupQuery only needs columns already in the index

Aggregation strategy follows a similar pattern: a GROUP BY on a small number of distinct groups is usually implemented as a hash aggregate (build one running total per group in memory), while a GROUP BY on data already sorted by the grouping key, often from a preceding sort or index scan, is implemented as a streaming aggregate that never needs the whole group set in memory at once.

Example

Running EXPLAIN ANALYZE SELECT * FROM users WHERE email = 'x@example.com' in PostgreSQL shows exactly which plan ran: an Index Scan using users_email_idx with actual row counts and timing means the index was used correctly; a Seq Scan on users means either no usable index exists or the optimizer judged one wasn’t worth it. MySQL’s equivalent, EXPLAIN FORMAT=JSON, and SQL Server’s graphical execution plan viewer expose the same underlying information: chosen access paths, join order, and estimated versus actual row counts.

ClickHouse and other OLAP-oriented engines optimize differently: instead of focusing heavily on join order for normalized schemas, they lean on columnar pruning (skipping whole data blocks using min/max statistics) and query-level parallelism across partitions, since analytical workloads are dominated by scan volume rather than join selectivity.

Dig deeper