Query Optimization and Execution Plans Questions
Making queries fast: reading and interpreting execution/explain plans, identifying full scans, spotting SQL anti-patterns, and rewriting queries for better performance. Covers how the planner chooses join order and access methods, and how statistics drive those choices. A core skill for anyone responsible for query performance in production.
Walk through the common physical operators you would see in a query execution plan (sequential scan, index scan, index-only scan, nested loop join, hash join, merge join, sort, aggregate). For each, explain why the optimizer would choose it and what cost trade-off (I/O vs. CPU vs. memory) it represents.
Sample Answer
Direct answer. Each physical operator represents a different strategy for reading or combining rows, and the optimizer's whole job is to pick a combination of these that minimizes total estimated cost for your specific query and data. Scans get rows out of storage; joins, sorts, and aggregates combine or reorganize the rows those scans produced.
Structured elaboration.
- A sequential scan reads every row of a table (or heap) in physical order. It's chosen when a large fraction of the table is needed, or when there's no useful index; its cost scales with table size regardless of how selective the filter is.
- An index scan walks an index structure to find matching rows, then fetches the full row from the table for each match. It's chosen when the predicate is selective enough that visiting the index plus a handful of table rows beats reading the whole table; the trade-off is that each match costs a separate (often random) I/O against the table.
- An index-only scan answers the query entirely from the index, with no visit to the table, when every column the query needs is present in the index and the storage engine's visibility bookkeeping allows it. This is the cheapest access method when it applies.
- Nested loop, hash, and merge join each combine two row sets differently: nested loop probes the inner side once per outer row (cheap for a small outer side with a cheap way to probe the inner side); hash join builds an in-memory hash table from one side and probes it with the other (good for large, unsorted inputs); merge join walks two already-sorted inputs in lockstep (good when both sides are cheaply available in sorted order).
- Sort materializes rows in a required order, which costs memory or disk depending on volume; aggregate collapses rows into groups, either via a hash table (unsorted input) or by exploiting already-sorted input.
Worked example. For SELECT customer_id, SUM(amount) FROM orders WHERE created_at >= '2025-01-01' GROUP BY customer_id, a reasonable plan is: an index scan on created_at to find recent rows cheaply (assuming that predicate is selective), feeding a hash aggregate that groups by customer_id using an in-memory hash table, since there's no reason to expect the rows already arrive sorted by customer.
Trade-offs and pitfalls. None of these operators is unconditionally "the good one" or "the bad one." A sequential scan is the CORRECT choice, not a mistake, when a query needs most of a table's rows, because the per-row overhead of index lookups would cost more in aggregate. The signal worth watching for in a plan is a mismatch between the operator and the actual selectivity or volume involved, for example a nested loop join running many more iterations than its own estimate expected.
Which planner configuration parameters most influence query plan choice for an analytical workload (think memory-per-operation settings and relative I/O cost settings)? For each, describe the direction of its effect on join selection, sort behavior, and scan-type choice, and how you would tune it safely in production rather than guessing.
Sample Answer
Direct answer. The settings that matter most for plan choice are the ones controlling how much memory a single operation (a sort or a hash table) is allowed before it spills to disk, and the ones controlling the RELATIVE cost the optimizer assigns to random versus sequential access; tune them by observing actual spill and scan behavior on representative queries, not by picking numbers from a generic guide.
Structured elaboration. A per-operation memory setting (commonly named something like "work memory") directly decides the threshold at which a sort or a hash join's build side spills to disk instead of staying in memory; set too low for the workload's typical operation sizes, and otherwise-reasonable plans start paying real spill costs, or the optimizer, anticipating that cost, avoids strategies (like a hash join) that would otherwise have been the best choice. A relative cost setting (often something like the assumed cost of a random-access page fetch versus a sequential one) directly shapes the scan-vs-index-scan decision: set it too high relative to the storage's real characteristics (a common legacy default tuned for spinning disks, not modern flash storage), and the optimizer systematically avoids index scans it should actually be choosing, favoring scans instead even when an index would genuinely be faster on the real hardware.
Tuning safely in production means changing one setting at a time, observing its effect on a representative sample of real queries (not just one synthetic test case), and rolling forward gradually rather than making a large jump across the whole fleet at once; because a memory setting is often a per-connection or per-query allowance, raising it globally multiplies total memory usage under concurrency, which needs to be weighed against the instance's actual available memory, not just its effect on one query in isolation.
Worked example. An instance still running with cost settings tuned for spinning-disk-era assumptions, now running entirely on flash storage, will systematically under-value index scans relative to their true cost on the current hardware; adjusting the relative random-versus-sequential cost setting downward, to better reflect flash storage's much smaller random-access penalty, can shift a whole class of borderline scan-vs-index-scan decisions toward the index scans that are now actually the better choice on the real hardware, without touching any individual query.
Trade-offs and pitfalls. These settings interact with each other and with the workload as a whole, changing one in isolation without observing its effect on the full mix of queries (not just the one you're trying to fix) risks helping one query while quietly hurting several others that were relying on the old defaults; validate with a representative before/after comparison across the workload, not a single query's before/after.
Compare IN, EXISTS, and JOIN as ways to test membership in another table. Cover how NULLs change the semantics of each (particularly for a large IN list or a NOT IN), when the optimizer is free to transform one into another, and when the choice actually changes performance rather than just readability.
Sample Answer
Direct answer. IN and EXISTS both test membership but differ in how NULLs interact with them and in how each maps onto a join in the optimizer's mind; JOIN differs from both in that it's row-producing rather than boolean, so it only behaves like a membership test if the inner side is guaranteed unique per outer row.
Structured elaboration. WHERE x IN (subquery) is true if x matches any row the subquery returns; critically, if the subquery's result set contains even one NULL and x doesn't match any non-NULL value, the whole expression evaluates to UNKNOWN rather than false, which under NOT IN specifically can silently make the entire outer predicate evaluate to nothing at all (a classic, expensive-to-debug correctness trap, not just a performance one). WHERE EXISTS (correlated subquery) instead asks only "does at least one matching row exist," is unaffected by NULLs in the same way, and lets the engine short-circuit as soon as one match is found rather than materializing a full list to compare against. A plain JOIN used for the same membership-testing purpose can silently change the row count of the outer query if the inner side isn't unique per join key, duplicating outer rows once per match, which neither IN nor EXISTS does since they only ever return true or false.
The optimizer is often free to transform one of these into another internally, when they happen to be logically equivalent for that specific query, so the choice between them is not always a performance decision. Where it genuinely matters, EXISTS's short-circuit behavior tends to help most when only a small fraction of outer rows actually have any match, and its immunity to the NULL trap makes it the generally safer default whenever the inner column can contain NULLs.
Worked example. For customers where you want those with at least one high-value order, IN (SELECT customer_id FROM orders WHERE amount > 1000) and EXISTS (SELECT 1 FROM orders o WHERE o.customer_id = c.customer_id AND o.amount > 1000) return identical results as long as customer_id in orders never contains NULL; if it can be NULL and you instead wrote NOT IN (...) to find customers WITHOUT a high-value order, that query can silently return zero rows the moment even one NULL customer_id exists in orders, while the equivalent NOT EXISTS form is unaffected.
Trade-offs and pitfalls. Prefer EXISTS/NOT EXISTS over IN/NOT IN whenever the inner column isn't guaranteed NOT NULL, purely for correctness, independent of any performance difference; and never reach for a plain JOIN as a membership test unless you've confirmed the inner side is unique per join key, or wrap it in a DISTINCT to restore the semantics you actually wanted.
Complexity
All three can, in principle, execute as a semi-join internally when the engine recognizes the pattern, so there is often no fundamental asymptotic difference; the practical differences come from NULL handling and from whether the engine actually recognizes the equivalence for a given query shape.
Edge cases
NOT IN against a column that can contain NULL is the single highest-risk pattern here and deserves an explicit check (or an automatic ban) in code review, since it fails silently rather than with an error.
A query filters on a column that has an index, but wrapping that column in a function or an implicit type conversion is silently preventing the index from being used. Walk through how you would confirm that is what's happening, and the different ways you could restore index usage (query-side and, where appropriate, schema-side).
Sample Answer
Direct answer. Confirm it by checking the WHERE clause literally for a function call or type cast wrapped around the indexed column; the fix is to rewrite the predicate so the indexed column appears bare on one side of the comparison, moving any transformation to the constant instead.
Structured elaboration. An index on a column can only be used efficiently by predicates the engine can translate directly into a range or equality scan of that index's stored values. The moment the column itself is wrapped in a function (a date-extraction function, a case-normalization function like lower()) or compared to a value of a different type that forces an implicit conversion, the engine generally can no longer map the predicate onto the index's stored order and has to fall back to evaluating the function for every row, which usually means a full scan. The confirmation step is mechanical: look at the WHERE clause and ask "is the indexed column, unmodified, on one side of a comparison operator?" If not, that's very likely the cause, and EXPLAIN will typically confirm it by showing a scan rather than the expected index usage.
The fix comes in two shapes. Query-side: rewrite the predicate so the column is bare and any transformation moves to the constant side of the comparison, for example turning an equality-on-a-truncated-date into a range comparison against the bare timestamp column. Schema-side, when a query-side rewrite genuinely isn't possible (the transformation is fundamental to the business logic, like case-insensitive matching), create an expression index that stores the transformed value directly, so the index itself already reflects lower(email) and the query can match against it.
Worked example. I verified the query-side rewrite is correctness-preserving with a small dataset: a query filtering CAST(order_ts AS DATE) = DATE '2025-01-01' (non-sargable, wraps the column) against four rows returns order_ids 1 and 3; rewriting to the equivalent range form returns the identical set.
-- non-sargable: function wraps the indexed column, defeats the index
SELECT order_id FROM orders
WHERE CAST(order_ts AS DATE) = DATE '2025-01-01';
-- sargable rewrite: bare column, range comparison against two constants
SELECT order_id FROM orders
WHERE order_ts >= TIMESTAMP '2025-01-01'
AND order_ts < TIMESTAMP '2025-01-02';
Both return order_id 1 and 3 for a table with rows at 2025-01-01 10:00, 2025-01-02 09:00, 2025-01-01 23:59:59, and 2025-02-01 00:00, confirming the rewrite changes only the execution strategy, not the result.
Trade-offs and pitfalls. The rewrite has to be exactly semantically equivalent, not just "close": an off-by-one on the upper bound (using <= against the start of the next day instead of < the next day) would silently include an extra midnight row. When an expression index is the only realistic fix, remember every query that wants to benefit from it must use the exact same expression the index was built on; a query written slightly differently (a different function, or the same function with different argument order) won't match.
Complexity
The rewrite doesn't change the query's asymptotic complexity by itself, it changes whether an O(log n) index lookup or an O(n) scan is even available as an option.
Edge cases
Time-zone-aware timestamp columns need extra care: a naive date-range rewrite can shift results by the UTC (Coordinated Universal Time) offset if the column and the literals aren't in the same time zone convention. Boundary values exactly at midnight need the range's inclusive/exclusive ends checked carefully against the original semantics.
An EXPLAIN ANALYZE shows a hash join spilling to disk (temp files). What causes a hash join to spill, how do you confirm that is actually happening from the plan output, and what are your options (query-level and configuration-level) for avoiding it?
Sample Answer
Direct answer. A hash join spills to disk when the hash table being built from one input doesn't fit inside the memory budget allotted to that operation, forcing the engine to partition both inputs and process them in batches with intermediate temp files; you confirm it from the plan by looking for an explicit batch or spill indicator alongside a jump in actual time relative to the row counts involved.
Structured elaboration. The build side of a hash join needs to fit (or be partitioned to fit) within the per-operation memory setting. When it doesn't, the engine splits both the build and probe inputs into multiple partitions small enough to fit in memory one at a time, writing the overflow to temporary disk files and reading them back in multiple passes. This isn't a bug: it's a graceful degradation that keeps the hash join correct at any input size, just at a real I/O and CPU cost.
To confirm this is happening (not, say, a slow input scan), most engines with detailed EXPLAIN ANALYZE output will explicitly report something like a batch count greater than one, or temp bytes/files written, directly on the hash join node; a hash join whose actual time is dramatically larger than a rough (rows times per-row cost) estimate, with no such explicit flag, is worth checking for spill via whatever temp-file or memory-usage view your engine exposes as a secondary confirmation.
Options to prevent or reduce spilling: increase the per-operation memory setting, but be conscious that this is a per-connection or per-query allowance, and raising it broadly can multiply total memory usage under concurrency; reduce the build side's row count with better filtering before the join (an earlier predicate or a narrower projection that lets more rows fit per memory unit); or reconsider whether the SMALLER of the two inputs is actually the one being chosen as the build side, since a misjudged build side is often the actual root cause rather than the memory setting itself.
Worked example. A feature-computation job whose hash join build side was under-estimated at 100,000 rows but turned out to be 5 million rows would very plausibly spill under a memory setting sized for the smaller estimate; fixing the underlying cardinality estimate (so the optimizer picks a bigger memory allocation, or a different join algorithm and build side entirely) often resolves this more durably than simply raising memory limits.
Trade-offs and pitfalls. Raising memory settings is the fastest fix but the least targeted: it helps this one query at the cost of every concurrently-running query on the instance potentially claiming more memory too, which can create its own resource-pressure problems under load. Prefer fixing the underlying cardinality estimate or the build-side size when that's the actual root cause.
Unlock Full Question Bank
Get access to all Query Optimization and Execution Plans interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.