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
WHEREfilters 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:
ordershas 50 million rows, an index exists oncustomer_id, and the query isSELECT * 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
regionstable against a 200-million-roworderstable onregion_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
regionsside, then streamordersthrough 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 BYor 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) orUPDATE 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 byWHERE region = 'EU'), andorder_items(600M rows). - Step: if the optimizer joins
orderstoorder_itemsfirst (both huge, no filter applied yet), it builds an enormous intermediate result before ever applying the selectivecustomersfilter. - Step: if instead it joins the small, filtered
customersset toordersfirst, cuttingordersdown to a few thousand matching rows, then joins that small intermediate result toorder_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:
EXPLAINalone shows what the optimizer expected;EXPLAIN ANALYZEshows 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 ANALYZEoutput 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 oncreated_atat 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 ANALYZEactually 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
WHEREclauses (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 commoncustomer_idwith 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 algorithm | Best when | Memory needs | Complexity |
|---|---|---|---|
| Nested Loop | One side is small, or an index exists on the join key | Minimal | Simple, O(n × m) worst case |
| Hash Join | Both sides large, unsorted, equality join | Needs to build a hash table for the smaller side | Moderate, roughly O(n + m) |
| Merge Join | Both sides already sorted on the join key | Minimal if pre-sorted, else needs a sort step | Efficient once sorted, O(n + m) |
| Scan type | Reads | Best when |
|---|---|---|
| Sequential Scan | Every row in the table | Filter matches most rows, or table is small |
| Index Scan | Index pages, then row lookups | Filter is selective, few matching rows |
| Index-Only Scan | Index pages only, no table lookup | Query 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.
Related Terms
- Database Indexing Internals — the access methods the optimizer is choosing between
- Database Normalization and Denormalization — schema shape drives how many joins the optimizer must plan around
- MVCC — visibility checks add per-row cost the optimizer must account for
- OLTP vs OLAP — the two workload types favor very different optimizer strategies
- Write-Ahead Logging (WAL) — write-heavy plans also account for logging cost, not just read cost
Referenced by