SQL courseLesson 18 of 23
SQL course · Lesson 18 of 23
SQL Query Optimization and Execution Plans
A practical process for speeding up slow SQL: read EXPLAIN plans, avoid full table scans, rewrite queries for performance and spot common anti-patterns.
On this page
- Sample data
- Query optimization and cost reduction
- What it is
- The process
- Pitfalls
- In interviews
- Reading EXPLAIN and execution plans
- What it is
- A first plan
- A plan with a join
- Common node types
- What to look for, in order
- Other engines
- Pitfalls
- In interviews
- Avoiding full table scans
- What it is
- Index the filter
- Keep filters sargable
- When a full scan is correct
- In columnar warehouses
- Pitfalls
- In interviews
- Query rewriting for performance
- What it is
- Correlated subquery to join and aggregate
- NOT IN to NOT EXISTS
- OFFSET pagination to keyset pagination
- Other standard rewrites
- Pitfalls
- In interviews
- Anti-pattern detection
- What it is
- Code-level anti-patterns
- Finding expensive queries in PostgreSQL
- In warehouses
- Pitfalls
- In interviews
- Practice questions
- Key takeaways
Query optimisation is a process, not a list of tricks: measure, find where the work goes, reduce it, and check the answer has not changed. Data Engineers need it because slow transformations delay every downstream table, and in cloud warehouses slow usually also means expensive. This lesson works through that process on a 300,000-row PostgreSQL table, reading real plans, and then maps the ideas to columnar warehouses.
Sample data
Generated data so the plans are realistic. It takes a few seconds to load.
CREATE TABLE customers (
customer_id INT PRIMARY KEY,
email TEXT NOT NULL,
country TEXT NOT NULL,
created_at DATE NOT NULL
);
INSERT INTO customers
SELECT g, 'user' || g || '@example.com',
(ARRAY['GB','IN','US','DE','FR'])[1 + g % 5],
DATE '2024-01-01' + (g % 700)
FROM generate_series(1, 20000) AS g;
CREATE TABLE orders (
order_id BIGINT PRIMARY KEY,
customer_id INT NOT NULL REFERENCES customers(customer_id),
status TEXT NOT NULL,
amount NUMERIC(10,2) NOT NULL,
created_at TIMESTAMP NOT NULL
);
INSERT INTO orders
SELECT g,
1 + (g::bigint * 7919) % 20000,
CASE WHEN g % 50 = 0 THEN 'refunded' WHEN g % 10 = 0 THEN 'pending' ELSE 'shipped' END,
((g::bigint * 37) % 50000) / 100.0,
TIMESTAMP '2025-01-01' + (g % 525600) * INTERVAL '1 minute'
FROM generate_series(1, 300000) AS g;
ANALYZE customers;
ANALYZE orders;
Each customer has 15 orders; 90% of orders are shipped. ANALYZE collects the statistics the planner uses for its estimates.
Query optimization and cost reduction
What it is
Optimising a query means getting the same result with less work: fewer rows and bytes read, less data moved between nodes, less sorting and hashing, and less waiting on locks or queues. “Cost” depends on the platform:
| Platform | What you pay for | What to reduce |
|---|---|---|
| PostgreSQL, MySQL, SQL Server | Fixed hardware; contention with other queries | Pages read, CPU, memory spills, lock time |
| Snowflake | Warehouse credits per second the warehouse runs | Runtime, partitions scanned, spilling, queueing |
| BigQuery (on-demand) | Bytes scanned | Columns read and partitions or clusters pruned |
| Spark / Databricks | Cluster time | Shuffle volume, skew, file count and size |
The process
- Reproduce and measure. Record runtime, rows returned and data scanned. Check whether it is slow every time or only under load (queueing and locks are not query problems).
- Read the plan and find the expensive nodes. Usually one or two operators account for most of the time.
- Read less data. Select only needed columns, filter early on columns that allow indexes or pruning, and aggregate before joining when the result is coarser.
- Do less work per row. Remove repeated subqueries, unnecessary sorts and
DISTINCTs, and fix join fan-out. - Fix the physical layer if the query is already sensible: indexes, partitioning, clustering, statistics, memory settings.
- Verify. Compare row counts and key aggregates with the original. A faster wrong query is a regression.
Pitfalls
- Tuning without a plan, or tuning the wrong query. Use workload statistics (below) to find the queries that consume the most total time, which is frequency times duration.
- Measuring on a warm cache once. Run several times and compare
BUFFERS(pages read from cache versus disk), not only milliseconds. - Optimising a query whose real problem is upstream: reprocessing all history every night when an incremental load would read 1% of it.
In interviews
“How would you optimise a slow query?” is answered best as the process above, with an example. Interviewers listen for measuring first and verifying last; candidates who jump straight to “add an index” miss both. Mention the platform’s cost model: bytes scanned in BigQuery, warehouse time in Snowflake.
Reading EXPLAIN and execution plans
What it is
SQL says what you want; the optimiser chooses how: which index, which join algorithm, in what order. The chosen strategy is the execution plan, a tree of operators. EXPLAIN shows the plan with estimates; EXPLAIN ANALYZE runs the query and adds what actually happened.
A first plan
SET max_parallel_workers_per_gather = 0;
EXPLAIN (ANALYZE, BUFFERS) SELECT order_id, amount FROM orders WHERE customer_id = 42;
Seq Scan on orders (cost=0.00..6250.00 rows=15 width=14) (actual time=0.904..69.897 rows=15 loops=1)
Filter: (customer_id = 42)
Rows Removed by Filter: 299985
Buffers: shared hit=2500
Planning Time: 0.034 ms
Execution Time: 69.931 ms
How to read each part:
- Node type:
Seq Scanreads the whole table. - cost=0.00..6250.00: estimated startup cost and total cost in arbitrary planner units (based on settings such as
seq_page_costandcpu_tuple_cost). Only useful for comparing plans of the same query. - rows=15: estimated rows this node outputs. width is the estimated average row size in bytes.
- actual time=0.904..69.897: milliseconds to the first row and to the last row, per loop.
- rows=15 loops=1: actual rows per loop and how many times the node ran. Total rows = rows x loops.
- Rows Removed by Filter: 299985: the scan read 300,000 rows to return 15. That ratio is the clearest sign of wasted work.
- Buffers: shared hit=2500: 8 kB pages found in PostgreSQL’s cache;
read=would be pages fetched from the operating system or disk.
A plan with a join
Plans are read from the innermost (most indented) nodes outwards; each node feeds its parent.
SET max_parallel_workers_per_gather = 0;
CREATE INDEX orders_customer_idx ON orders (customer_id);
CREATE INDEX orders_created_idx ON orders (created_at);
EXPLAIN (ANALYZE, BUFFERS)
SELECT c.country, COUNT(*) AS orders, SUM(o.amount) AS revenue
FROM orders o
JOIN customers c ON c.customer_id = o.customer_id
WHERE o.created_at >= DATE '2025-06-01' AND o.created_at < DATE '2025-07-01'
GROUP BY c.country;
HashAggregate (cost=2740.28..2740.34 rows=5 width=43) (actual time=79.820..79.826 rows=5 loops=1)
Group Key: c.country
Batches: 1 Memory Usage: 24kB
-> Hash Join (cost=607.42..2417.91 rows=42982 width=9) (actual time=25.530..69.470 rows=43200 loops=1)
Hash Cond: (o.customer_id = c.customer_id)
-> Index Scan using orders_created_idx on orders o (cost=0.42..1698.06 rows=42982 width=10)
(actual time=0.071..12.960 rows=43200 loops=1)
Index Cond: ((created_at >= '2025-06-01'::date) AND (created_at < '2025-07-01'::date))
-> Hash (cost=357.00..357.00 rows=20000 width=7) (actual time=25.129..25.131 rows=20000 loops=1)
Buckets: 32768 Batches: 1 Memory Usage: 1038kB
-> Seq Scan on customers c (cost=0.00..357.00 rows=20000 width=7) (actual time=0.012..11.302 rows=20000 loops=1)
Execution Time: 80.113 ms
(Trimmed: buffer lines removed.) The story: read June’s orders through the created_at index, build a hash table of all customers, probe it for each order, then hash-aggregate into 5 countries. Estimates (42,982 rows) are close to actuals (43,200), so the planner had good information. Batches: 1 on the hash and aggregate means everything fit in memory.
Common node types
| Node | Meaning | Watch for |
|---|---|---|
| Seq Scan | Read every row | High Rows Removed by Filter on a selective query |
| Index Scan / Index Only Scan | Navigate an index, then fetch rows (or not) | Many loops; Heap Fetches on index-only scans |
| Bitmap Index Scan + Bitmap Heap Scan | Collect matching row locations, then read pages in order | lossy heap blocks when memory is short |
| Nested Loop | For each outer row, look up inner rows | Large loops with no index on the inner side |
| Hash Join | Build a hash table on one input, probe with the other | Batches > 1 (spilled to disk) |
| Merge Join | Merge two inputs sorted on the join key | Expensive sorts feeding it |
| Sort | Order rows | Sort Method: external merge Disk: means a spill |
| HashAggregate / GroupAggregate | GROUP BY | Disk Usage or several batches |
| Gather | Collect results from parallel workers | Workers planned but not launched |
What to look for, in order
- The node with the most time. Actual time is cumulative: a parent’s time includes its children’s. Subtract to find where the time is really spent.
- Estimate versus actual rows. A difference of 10x or more means the optimiser is guessing badly (stale or missing statistics, correlated columns, functions on columns), and the rest of the plan may be built on that bad guess.
- Loops. A cheap inner node that runs 100,000 times is not cheap.
- Spills. Disk usage in sorts, hashes and aggregates.
- Wasted reads. Rows removed by filter, or a scan returning far more rows than the final result needs.
Other engines
-- Snowflake: text plan, plus the graphical Query Profile in Snowsight
EXPLAIN USING TEXT SELECT ...;
-- Watch "Partitions scanned" versus "Partitions total", "Bytes spilled to local/remote storage",
-- and exploding joins in the Query Profile.
-- BigQuery: no EXPLAIN statement; the console's "Execution details" tab shows stages,
-- slot time, bytes shuffled and records read/written per stage.
-- A dry run estimates bytes processed before running: bq query --dry_run '...'
-- Spark SQL / Databricks
EXPLAIN FORMATTED SELECT ...;
Pitfalls
- Comparing
costnumbers across different queries or databases. They are relative units for one planner decision. - Forgetting that
actual timeandrowsare per loop. - Reading only the top line. The root shows the total; the problem is usually deep in the tree.
- Testing on a tiny development table. The planner correctly chooses a sequential scan on small tables; you need production-like volume and statistics.
In interviews
You may be shown a plan and asked what is wrong. Talk through it bottom-up: what is scanned, how it is joined, where the estimates diverge, where time goes and whether anything spills. Naming the three join algorithms and when each is chosen is a common follow-up (see the optimizer internals lesson in this course).
Avoiding full table scans
What it is
A full table scan reads every row (or every micro-partition, or every file). It is the right choice when the query needs a large share of the table, and the wrong one when it needs a handful of rows. Avoiding it means giving the engine a way to skip data: an index, partition pruning, or min/max metadata, and writing filters it can use.
Index the filter
The first plan above read 300,000 rows to find 15. With an index on customer_id:
SET max_parallel_workers_per_gather = 0;
EXPLAIN (ANALYZE, BUFFERS) SELECT order_id, amount FROM orders WHERE customer_id = 42;
Bitmap Heap Scan on orders (cost=4.54..61.24 rows=15 width=14) (actual time=0.056..0.088 rows=15 loops=1)
Recheck Cond: (customer_id = 42)
Heap Blocks: exact=15
Buffers: shared hit=15 read=3
-> Bitmap Index Scan on orders_customer_idx (cost=0.00..4.54 rows=15 width=0) (actual time=0.044..0.045 rows=15 loops=1)
Index Cond: (customer_id = 42)
Buffers: shared read=3
Execution Time: 0.170 ms
18 pages instead of 2,500, and well under a millisecond instead of about 70 ms on this machine.
Keep filters sargable
A predicate is sargable (Search ARGument ABLE) when the engine can match it against an index or partition metadata. Applying a function or cast to the column usually breaks that, because the index stores the raw column values.
SET max_parallel_workers_per_gather = 0;
EXPLAIN (ANALYZE) SELECT count(*) FROM orders WHERE created_at::date = DATE '2025-03-01';
EXPLAIN (ANALYZE) SELECT count(*) FROM orders
WHERE created_at >= DATE '2025-03-01' AND created_at < DATE '2025-03-02';
-- cast on the column: full scan
Aggregate (actual time=115.691..115.693 rows=1 loops=1)
-> Seq Scan on orders (cost=0.00..7000.00 rows=1500 width=0) (actual time=50.090..115.592 rows=1440 loops=1)
Filter: ((created_at)::date = '2025-03-01'::date)
Rows Removed by Filter: 298560
-- half-open range on the raw column: index range scan
Aggregate (actual time=0.534..0.534 rows=1 loops=1)
-> Index Only Scan using orders_created_idx on orders (cost=0.42..67.12 rows=1535 width=0) (actual time=0.114..0.474 rows=1440 loops=1)
Index Cond: ((created_at >= '2025-03-01'::date) AND (created_at < '2025-03-02'::date))
Same 1,440 rows; the second form uses the index. Other common non-sargable forms and their fixes:
| Non-sargable | Sargable alternative |
|---|---|
WHERE YEAR(order_date) = 2025 |
WHERE order_date >= '2025-01-01' AND order_date < '2026-01-01' |
WHERE amount * 1.2 > 100 |
WHERE amount > 100 / 1.2 |
WHERE lower(email) = '...' |
an expression index on lower(email), or store normalised values |
WHERE customer_id::text = '42' |
compare with the column’s own type: customer_id = 42 |
WHERE email LIKE '%@example.com' |
reverse-string index, trigram index, or a separate domain column |
WHERE COALESCE(region, 'x') = 'EU' |
WHERE region = 'EU' (NULLs never equal ‘EU’ anyway) |
An expression index makes a function call sargable when you cannot change the query:
SET max_parallel_workers_per_gather = 0;
CREATE INDEX customers_email_lower_idx ON customers (lower(email));
EXPLAIN SELECT * FROM customers WHERE lower(email) = 'user42@example.com';
Bitmap Heap Scan on customers (cost=5.06..151.93 rows=100 width=32)
Recheck Cond: (lower(email) = 'user42@example.com'::text)
-> Bitmap Index Scan on customers_email_lower_idx (cost=0.00..5.04 rows=100 width=0)
Index Cond: (lower(email) = 'user42@example.com'::text)
The rows=100 estimate is a default guess (0.5% of the table) because the new expression has no statistics yet; running ANALYZE customers collects them.
When a full scan is correct
SET max_parallel_workers_per_gather = 0;
EXPLAIN SELECT order_id FROM orders WHERE status = 'shipped';
Seq Scan on orders (cost=0.00..6250.00 rows=270010 width=8)
Filter: (status = 'shipped'::text)
90% of rows match, so reading the table sequentially is cheaper than jumping through an index for each row. An index only helps a selective filter; the crossover point depends on row width, caching and storage, which is why the planner decides with statistics rather than a fixed rule.
In columnar warehouses
Snowflake, BigQuery, Redshift and Delta Lake have no B-tree indexes for analytical tables. They skip data with pruning: each micro-partition, file or row group records min/max values per column, and the engine skips those whose range cannot match. The same sargability rules apply. A filter on a partition or cluster column that wraps the column in a function can disable pruning, and pruning only helps when the data is physically organised by that column (partitioning, clustering keys, Z-ordering). Selecting fewer columns also reduces the scan, because columnar storage reads each column separately.
Pitfalls
- Indexing every column. Each index slows writes and uses space; index for the queries you actually run.
- Implicit casts from mismatched join or filter types (an
INTcolumn compared with a text parameter from an application) silently cause full scans. ORacross different columns can still use indexes in PostgreSQL (aBitmapOr), but many engines fall back to a scan; aUNION ALLof two selective queries is the usual rewrite.- Stale statistics make the planner pick a scan (or an index) wrongly. Check when tables were last analysed.
In interviews
Expect “this query on a date column is slow; why?” with a CAST or DATE() in the WHERE clause. Explain sargability, give the range rewrite, and mention that in a warehouse the same mistake prevents partition pruning. Saying when a full scan is the right plan shows you understand selectivity rather than following rules.
Query rewriting for performance
What it is
Two queries can return the same rows with very different plans. Modern optimisers rewrite many forms themselves (for example, PostgreSQL turns IN (subquery) into a semi-join), but not all of them, and not in every engine. Knowing the standard rewrites lets you help the optimiser, and lets you recognise when it already did the work.
Correlated subquery to join and aggregate
SET max_parallel_workers_per_gather = 0; SET jit = off;
EXPLAIN (ANALYZE)
SELECT c.customer_id,
(SELECT SUM(o.amount) FROM orders o WHERE o.customer_id = c.customer_id) AS total
FROM customers c
WHERE c.country = 'GB';
Seq Scan on customers c (cost=0.00..245569.26 rows=4000 width=36) (actual time=0.734..150.045 rows=4000 loops=1)
Filter: (country = 'GB'::text)
SubPlan 1
-> Aggregate (cost=61.28..61.29 rows=1 width=32) (actual time=0.035..0.035 rows=1 loops=4000)
-> Bitmap Heap Scan on orders o (actual time=0.008..0.030 rows=15 loops=4000)
-> Bitmap Index Scan on orders_customer_idx (actual time=0.004..0.004 rows=15 loops=4000)
Execution Time: 150.421 ms
The subquery runs once per customer (loops=4000). The set-based rewrite:
SET max_parallel_workers_per_gather = 0; SET jit = off;
EXPLAIN (ANALYZE)
SELECT c.customer_id, SUM(o.amount) AS total
FROM customers c
LEFT JOIN orders o ON o.customer_id = c.customer_id
WHERE c.country = 'GB'
GROUP BY c.customer_id;
HashAggregate (cost=7044.67..7094.67 rows=4000 width=36) (actual time=154.012..155.902 rows=4000 loops=1)
Group Key: c.customer_id
-> Hash Right Join (cost=457.00..6744.67 rows=60000 width=10) (actual time=2.722..123.374 rows=60000 loops=1)
Hash Cond: (o.customer_id = c.customer_id)
-> Seq Scan on orders o (actual time=0.008..62.847 rows=300000 loops=1)
-> Hash (actual time=2.699..2.700 rows=4000 loops=1)
-> Seq Scan on customers c (actual time=0.010..2.208 rows=4000 loops=1)
Execution Time: 156.318 ms
On this data both take about the same time: the correlated version is saved by the index, which makes each of the 4,000 lookups cheap. The difference shows up as data grows or when the index is missing: without orders_customer_idx, the correlated plan’s estimated cost rises to about 25 million units (a full scan of orders per customer), while the join plan does not change. The lesson is to measure: rewrites are hypotheses, and the plan tells you whether they helped. In distributed engines, which often cannot run a correlated subquery as a loop at all, the join-and-aggregate form is the standard.
The LEFT JOIN keeps customers with no orders (their total is NULL, as in the scalar subquery). An inner join would silently drop them.
NOT IN to NOT EXISTS
NOT IN against a subquery that can return NULL gives a different answer, not just a slower one:
CREATE TABLE blocked (customer_id INT);
INSERT INTO blocked VALUES (1), (2), (NULL);
SELECT count(*) AS not_in FROM customers
WHERE customer_id NOT IN (SELECT customer_id FROM blocked);
SELECT count(*) AS not_exists FROM customers c
WHERE NOT EXISTS (SELECT 1 FROM blocked b WHERE b.customer_id = c.customer_id);
DROP TABLE blocked;
not_in is 0 and not_exists is 19998. x NOT IN (1, 2, NULL) is never true, because x <> NULL is unknown. NOT EXISTS has the intended meaning and is planned as an anti-join; prefer it (or LEFT JOIN ... WHERE right.key IS NULL).
OFFSET pagination to keyset pagination
SET max_parallel_workers_per_gather = 0;
EXPLAIN (ANALYZE) SELECT order_id, amount FROM orders ORDER BY order_id LIMIT 20 OFFSET 250000;
EXPLAIN (ANALYZE) SELECT order_id, amount FROM orders WHERE order_id > 250000 ORDER BY order_id LIMIT 20;
-- OFFSET: reads and discards 250,000 rows
Limit (actual time=129.856..129.861 rows=20 loops=1)
-> Index Scan using orders_pkey on orders (actual time=0.030..117.632 rows=250020 loops=1)
Execution Time: 129.878 ms
-- keyset: seeks straight to the position
Limit (actual time=0.042..0.048 rows=20 loops=1)
-> Index Scan using orders_pkey on orders (actual time=0.041..0.045 rows=20 loops=1)
Index Cond: (order_id > 250000)
Execution Time: 0.060 ms
OFFSET n still produces the first n rows and throws them away, so later pages get slower. Keyset (seek) pagination remembers the last key returned and filters on it. The same idea makes batch extracts restartable: WHERE id > :last_id ORDER BY id LIMIT 10000.
Other standard rewrites
| Pattern | Rewrite | Why |
|---|---|---|
Join, then GROUP BY at a coarser grain |
Aggregate the big table first, then join | Joins fewer rows |
SELECT DISTINCT to remove join duplicates |
Fix the join key, or use EXISTS |
DISTINCT hides fan-out and adds a sort or hash |
UNION when duplicates are impossible |
UNION ALL |
Skips a de-duplication step |
COUNT(*) > 0 to test existence |
EXISTS (...) |
Stops at the first match |
| Same CTE or subquery computed several times | Compute once (temporary table or materialised CTE) | Avoids repeated work |
OR across columns with separate indexes |
UNION ALL of two selective queries (guarding duplicates) |
Lets each branch use its index |
| Row-by-row loop in application code | One set-based statement | Removes per-row round trips |
Pitfalls
- Changing the result.
LEFTtoINNER,NOT INtoNOT EXISTSwith NULLs,UNIONtoUNION ALLwith real duplicates: each can change the answer. Compare counts and checksums before and after. - Assuming the optimiser does not already do it. Check the plan of the original first.
- Keyset pagination needs a unique, indexed order. Use a tiebreaker such as
(created_at, id) > (:ts, :id).
In interviews
Expect to be given a correlated subquery or NOT IN and asked to improve it. Rewrite it, explain why it is faster, and explicitly state how you would prove the results are identical. Bringing up the NOT IN NULL trap unprompted is a strong signal.
Anti-pattern detection
What it is
An anti-pattern is a query shape that is usually slow or wrong. Detecting them has two halves: reviewing code for known shapes, and finding expensive queries in production from workload statistics, then reading their plans.
Code-level anti-patterns
| Anti-pattern | Why it hurts | Fix |
|---|---|---|
SELECT * in pipelines |
Reads every column (costly in columnar engines); breaks when the schema changes | List columns |
| Function or cast on a filtered column | Disables indexes and pruning | Range predicates, expression indexes |
Leading wildcard LIKE '%x' |
Cannot use a B-tree index | Trigram or full-text index, derived column |
NOT IN (subquery) |
Wrong with NULLs | NOT EXISTS |
DISTINCT to hide duplicates |
Masks a join bug, adds work | Fix the join grain |
| Correlated subquery per row | Runs once per outer row | Join and aggregate |
OFFSET deep pagination |
Reads and discards rows | Keyset pagination |
| Join on mismatched types | Implicit cast, no index use | Align types in the model |
| Unbounded cross join or missing join condition | Row explosion | Explicit JOIN ... ON with keys |
ORDER BY in subqueries and views |
Wasted sort; order is not guaranteed outside | Sort once at the end |
Huge IN (...) lists built in code |
Long parse and plan, plan-cache misses | Join to a temporary table or array parameter |
Linters such as SQLFluff catch some of these (for example SELECT * and implicit joins) in CI before code reaches production.
Finding expensive queries in PostgreSQL
The pg_stat_statements extension aggregates statistics per normalised query. It must be loaded through shared_preload_libraries, so it is shown here without running:
-- Requires shared_preload_libraries = 'pg_stat_statements' and a restart
CREATE EXTENSION IF NOT EXISTS pg_stat_statements;
SELECT left(query, 60) AS query,
calls,
round(total_exec_time) AS total_ms,
round(mean_exec_time, 2) AS mean_ms,
rows,
shared_blks_read + shared_blks_hit AS pages
FROM pg_stat_statements
ORDER BY total_exec_time DESC
LIMIT 10;
Sort by total time first: a 20 ms query called a million times matters more than a 10-second report run once a day. Two other tools help: log_min_duration_statement logs every statement slower than a threshold, and the auto_explain module logs the plans of slow statements automatically, so you see the plan that actually ran.
The table statistics views reveal tables that are scanned in full far more often than through indexes:
SELECT relname, seq_scan, seq_tup_read, idx_scan
FROM pg_stat_user_tables
WHERE relname IN ('orders', 'customers')
ORDER BY relname;
A large table with a high seq_tup_read and few idx_scans is a candidate for an index or a query fix. (On this freshly loaded sample the numbers just reflect the examples above.)
In warehouses
-- Snowflake: heaviest queries in the last day from ACCOUNT_USAGE (data can lag behind real time)
SELECT query_id, LEFT(query_text, 80) AS query, total_elapsed_time / 1000 AS seconds,
partitions_scanned, partitions_total, bytes_spilled_to_local_storage
FROM snowflake.account_usage.query_history
WHERE start_time > DATEADD('day', -1, CURRENT_TIMESTAMP())
ORDER BY total_elapsed_time DESC
LIMIT 10;
-- BigQuery: most bytes billed in the last day
SELECT job_id, LEFT(query, 80) AS query, total_bytes_billed, total_slot_ms
FROM `region-us`.INFORMATION_SCHEMA.JOBS_BY_PROJECT
WHERE creation_time > TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 1 DAY)
ORDER BY total_bytes_billed DESC
LIMIT 10;
Signals to look for: partitions scanned close to partitions total (no pruning), bytes spilled to remote storage (warehouse too small or a huge sort or join), join output much larger than its inputs (exploding join), and long queue time (concurrency, not query shape).
Pitfalls
- Fixing the slowest single query rather than the biggest total consumer.
- Trusting
EXPLAINwithoutANALYZEwhen estimates are wrong; the problem is often the misestimate itself. - Treating every sequential scan as an anti-pattern. Small tables and low-selectivity filters should be scanned.
In interviews
“How would you find the queries that are costing us the most?” tests whether you know the tooling: pg_stat_statements and slow-query logs in PostgreSQL, QUERY_HISTORY and Query Profile in Snowflake, INFORMATION_SCHEMA.JOBS in BigQuery. Pair that with a short list of anti-patterns you would look for in the top offenders and how you would prevent them in code review (linters, review checklists, cost monitoring).
Practice questions
A query filters with WHERE DATE(created_at) = ‘2025-03-01’ and is slow despite an index on created_at. Why, and how do you fix it?
The function on the column makes the predicate non-sargable: the index stores created_at values, not DATE(created_at), so the engine scans every row. Rewrite as a half-open range: created_at >= '2025-03-01' AND created_at < '2025-03-02'. Alternatively, create an expression index on the date, but the range rewrite is the general fix and also enables partition pruning in warehouses.
In EXPLAIN ANALYZE, a node shows rows=1 estimated and rows=250000 actual. What does that tell you?
The optimiser badly underestimated the cardinality, so it probably chose a plan suited to one row (such as a nested loop with an index lookup per row) that is terrible for 250,000. Causes include stale statistics, correlated predicates, functions on columns and skewed values. Run ANALYZE, check the predicate, and consider extended statistics; then re-read the plan.
Why can NOT IN return no rows when you expect many?
If the subquery returns any NULL, x NOT IN (..., NULL) evaluates to unknown for every x that is not in the list, and unknown is not true, so no rows pass. Use NOT EXISTS (or filter NULLs out of the subquery).
Your nightly export uses LIMIT 10000 OFFSET n and gets slower every page. What do you change?
Switch to keyset pagination: order by a unique indexed key and filter on the last key seen (WHERE id > :last_id ORDER BY id LIMIT 10000). Each page then seeks directly instead of reading and discarding all earlier rows, and the export can resume from the last key after a failure.
When is a full table scan the best plan?
When the query needs a large fraction of the rows (low selectivity), when the table is small, or for aggregations over the whole table. Sequential reads are much cheaper per row than random index lookups, so beyond some fraction of the table the scan wins. In columnar engines you still reduce work by reading only needed columns.
How do you find which queries to optimise first in PostgreSQL?
Use pg_stat_statements and sort by total execution time (calls times mean), then look at mean time, rows and buffer usage. Get the real plans with auto_explain or by running EXPLAIN (ANALYZE, BUFFERS) on representative parameters. Fix the top consumers first and confirm the improvement in the same statistics afterwards.
You rewrote a query and it is now 10x faster. What do you check before deploying?
That it returns the same result: row counts, sums or checksums of key columns, and a diff of the outputs (EXCEPT in both directions) on representative data, including edge cases such as NULLs, customers without orders and duplicates. Also check the plan on production-sized data, not just a development sample.
Key takeaways
- Optimise as a process: measure, read the plan, read less data, do less work per row, fix the physical layer, verify.
- In a plan, find the node with the most time, compare estimated with actual rows, and watch loops, spills and rows removed by filters.
- Keep filters sargable: no functions or casts on filtered or joined columns; use ranges and matching types.
- A full scan is correct for low-selectivity queries; indexes and pruning help only selective filters.
- Standard rewrites (join and aggregate instead of correlated subqueries,
NOT EXISTSinstead ofNOT IN, keyset pagination) are hypotheses: confirm them with the plan and the results. - Find what to fix from workload statistics (
pg_stat_statements, query history), sorted by total time or cost.
Progress is saved in this browser only. No account needed.