Query Optimization: Transforming Declarative SQL into Optimal Physical Execution Plans

SQL is a purely declarative language: a developer specifies *what* data to retrieve rather than *how* to retrieve it algorithmically. A simple three-table join can be executed in thousands of mathematically equivalent ways, with execution runtimes ranging from milliseconds to hours depending on index usage, join algorithms, and access order.

The **Query Optimizer** is the computational brain of a database management system. It translates parsed declarative SQL Abstract Syntax Trees (ASTs) into an optimal physical execution plan by applying relational algebra equivalence rules, estimating intermediate result cardinalities, and evaluating total I/O and CPU cost across the search space.

The Query Compilation & Optimization Pipeline

A declarative query progresses through four distinct architectural stages before physical execution:

  1. Parsing & Semantic Analysis: Validates SQL syntax against system catalog schemas (table existence, column types, permissions) to produce an initial AST.
  2. Logical Optimization (Rule-Based Rewrites): Applies heuristic relational algebra rewrites (such as Predicate Pushdown, Projection Pruning, and Subquery Flattening) to produce a normalized Logical Plan.
  3. Cardinality & Cost Estimation: Evaluates statistical histograms, HyperLogLog distinct estimates, and single/multi-column correlation factors to calculate the expected row output for every sub-operator.
  4. Physical Plan Generation: Translates logical operators into concrete algorithmic implementations (e.g., converting a Logical Join into a Hash Join, Merge Join, or Index Nested Loop Join) and selects the lowest-cost candidate plan.

Logical Transformations via Relational Algebra

Optimizers apply equivalence rules to transform query trees into smaller intermediate shapes prior to physical evaluation:

  • Predicate Pushdown: Moves filter conditions (Selection $\sigma$) as deep down the operator tree as possible, filtering rows directly during disk page scans before expensive joins occur.
  • Projection Pruning: Eliminates unused columns (Projection $\pi$) early in the plan pipeline to minimize memory footprint in join buffers.
  • Subquery Unnesting / Decorrelation: Converts correlated subqueries and `EXISTS` clauses into semi-joins or inner joins, avoiding quadratic $O(N \times M)$ nested loop iterations.

Join Ordering: Dynamic Programming vs. Top-Down Frameworks

For a query joining $N$ tables, the number of possible join tree topologies grows factorially according to Catalan numbers ($O(\frac{(2N-2)!}{(N-1)!})$), making exhaustive brute-force search impossible for large $N$.

1. Bottom-Up Dynamic Programming (System R Style)

Pioneered by Patricia Selinger in 1979 for IBM System R. It builds optimal join trees incrementally from sub-plans of size $k$ up to size $N$:

  • Finds optimal single-table access paths (Seq Scan vs. Index Scan).
  • Combines pairs of relations into optimal 2-way joins, pruning strictly dominated plans.
  • Iteratively extends to $(k+1)$-way joins using memoized optimal $k$-way results.
  • Tracks **Interesting Orders** (preserving sorted physical output needed for downstream `ORDER BY` or `GROUP BY` to avoid redundant sorting).

2. Top-Down Extensible Search (Volcano / Cascades Framework)

Designed by Goetz Graefe, the **Cascades Framework** organizes candidate plans into equivalence classes (Groups). It explores logical transformation rules and physical implementation rules using goal-driven, memoized branch-and-bound pruning (used by Microsoft SQL Server, CockroachDB, and Apache Calcite).

Physical Join Implementations

The Cost-Based Optimizer maps logical joins to one of three primary physical join algorithms based on estimated cardinalities and available indexes:

  • Nested Loop Join / Index Nested Loop: Iterates outer rows and queries inner rows. Optimal when the outer set is small ($< 1,000$ rows) and the inner relation possesses a selective B+ Tree index ($O(N \log M)$).
  • Hash Join: Builds an in-memory hash table on the smaller relation (Build phase), then streams and probes the larger relation (Probe phase). Optimal for large unsorted datasets without usable indexes ($O(N + M)$).
  • Sort-Merge Join: Sorts both inputs by join key (if not already sorted by an index scan) and advances parallel iterators through matching rows. Optimal for massive data volumes where data is pre-sorted or memory is limited.

C++ Conceptual Simulation Blueprint (Dynamic Programming Join Ordering)

#include <iostream>
#include <vector>
#include <string>
#include <unordered_map>
#include <algorithm>
#include <climits>

struct Plan {
    std::string description;
    double cost;
    uint64_t cardinality;
};

class DPJoinOptimizer {
private:
    // Bitmask memoization table: key is set of relations (e.g., bitmask 0b111 = {R1, R2, R3})
    std::unordered_map<uint32_t, Plan> memo;
    std::vector<std::string> tableNames;
    std::vector<uint64_t> tableSizes;

public:
    DPJoinOptimizer(std::vector<std::string> names, std::vector<uint64_t> sizes)
        : tableNames(names), tableSizes(sizes) {}

    Plan findOptimalPlan() {
        int n = tableNames.size();
        
        // 1. Base case: single-table access paths (size 1)
        for (int i = 0; i < n; ++i) {
            uint32_t mask = (1 << i);
            memo[mask] = {tableNames[i], static_cast<double>(tableSizes[i]), tableSizes[i]};
        }

        // 2. Iterate subset sizes from 2 to N
        for (int size = 2; size <= n; ++size) {
            for (uint32_t mask = 1; mask < (1 << n); ++mask) {
                if (__builtin_popcount(mask) != size) continue;

                Plan bestPlan = {"", 1e18, 0};

                // Partition mask into two non-empty disjoint submasks: S1 and S2
                for (uint32_t s1 = (mask - 1) & mask; s1 > 0; s1 = (s1 - 1) & mask) {
                    uint32_t s2 = mask ^ s1;
                    if (s1 > s2) continue; // Symmetric avoidance

                    if (memo.count(s1) && memo.count(s2)) {
                        const auto& p1 = memo[s1];
                        const auto& p2 = memo[s2];

                        // Simplified Cost Model: Hash Join cost = P1_card + P2_card + join_result
                        uint64_t joinedCard = (p1.cardinality * p2.cardinality) / 1000; // Selectivity factor
                        double joinCost = p1.cost + p2.cost + (p1.cardinality + p2.cardinality);

                        if (joinCost < bestPlan.cost) {
                            bestPlan = {"(" + p1.description + " JOIN " + p2.description + ")", joinCost, joinedCard};
                        }
                    }
                }
                memo[mask] = bestPlan;
            }
        }
        return memo[(1 << n) - 1];
    }
};

Real-World Database Systems & Optimizer Architectures

  1. PostgreSQL Query Planner: Employs bottom-up dynamic programming (standard System R) for queries with fewer than 12 joins, switching to Genetic Query Optimization (GEQO) for larger joins to bound planning latency.
  2. CockroachDB & Microsoft SQL Server: Implements top-down Cascades frameworks to seamlessly optimize complex distributed and partitioned SQL transformations with extensible rule engines.
  3. Apache Calcite: Extensible query optimization framework powering SQL compilation and physical rule selection for systems like Apache Flink, Apache Hive, and Drill.
  4. Presto / Trino Distributed Cost Optimizer: Dynamically plans distributed shuffle exchanges, broadcast hash joins, and partition repartitioning across worker nodes based on table statistics.