SQL query optimization is one of those topics where generic advice is almost useless. "Filter early, join less, avoid SELECT *" is correct but incomplete. The specific techniques that make a meaningful difference depend on the query engine, the storage format, the distribution of data, and the shape of the query itself.
This article covers SQL optimization techniques that apply specifically to columnar analytics warehouses (Redshift, BigQuery, Snowflake, DuckDB, Athena), explaining why each technique matters in the columnar execution model and where each one is most effective.
Column projection: read what you need
In a columnar store, the query planner reads only the columns referenced in the query. SELECT * forces the engine to read every column from every row group, eliminating most of the I/O advantage of columnar storage. Always select specific columns, even in exploratory queries.
This extends to CTEs and subqueries. A CTE that selects twenty columns from a large table and is then consumed by an outer query that uses only three of those columns may or may not benefit from column pruning, depending on whether your warehouse's optimizer can push the projection into the CTE. Snowflake and BigQuery typically do. Redshift sometimes does not for complex CTEs. When in doubt, explicitly limit the columns in the CTE to those the outer query uses.
Predicate placement: where filter conditions execute
Filters reduce the data that needs to be processed. The earlier in the execution plan a filter applies, the more downstream work it eliminates. The optimizer should push filters down to the scan layer wherever possible, but optimizer decisions are not always correct and SQL constructs can block pushdown.
Filter conditions on partition columns should be written as direct comparisons to literal values or simple expressions, not as function calls on the column. WHERE order_date >= '2025-01-01' enables partition pruning. WHERE DATE_TRUNC('month', order_date) = '2025-01-01' applies a function to the column before comparison, which prevents the optimizer from matching the filter to the partition range. Both queries return the same result, but the first scans one or two partitions while the second may scan all partitions.
In views and subqueries, verify that filters from the outer query are pushed into the inner scan. A view defined as SELECT * FROM events WHERE event_type = 'purchase' combines with an outer WHERE user_id = 123; the question is whether the query engine pushes the user_id filter into the underlying table scan or applies it after reading all purchase events. Check the query plan to confirm.
Join ordering and the cost of large-to-large joins
In distributed execution, join order matters because it determines how much data is shuffled across the network. The general principle is to reduce data size early: filter before joining, join smaller tables first, and structure joins so the data flowing into each successive join is smaller than the data flowing in from the previous one.
Query optimizers estimate join ordering costs using table statistics. When the optimizer has accurate statistics for all tables in the query, it often chooses a good order. When statistics are stale, it may choose a poor order that processes large intermediate results unnecessarily.
One specific join anti-pattern: joining two large tables without filtering either one first. A fact table with 500 million rows joined against a dimension table with 10 million rows (both large) is significantly more expensive than a fact table filtered to the relevant date range (perhaps 50 million rows) joined against the same dimension table. The filter does not change the logical result, but it changes the cost of the join substantially. Write the query to filter before the join, not after.
Aggregation before join
A frequent optimization opportunity is reordering operations so aggregation happens before a join rather than after. Consider a query that computes revenue by region by joining orders against a region lookup table and summing:
-- Common but expensive pattern
SELECT r.region_name, SUM(o.amount) AS total_revenue
FROM orders o
JOIN regions r ON o.region_code = r.region_code
WHERE o.order_date >= '2025-01-01'
GROUP BY r.region_name
-- More efficient: aggregate first, then join
SELECT r.region_name, agg.total_revenue
FROM (
SELECT region_code, SUM(amount) AS total_revenue
FROM orders
WHERE order_date >= '2025-01-01'
GROUP BY region_code
) agg
JOIN regions r ON agg.region_code = r.region_code
The second form aggregates millions of order rows into at most a few thousand region-level totals before joining against the regions lookup table. The first form joins millions of rows against the regions table, then aggregates. Both produce the same result. The second form moves significantly less data through the join.
Not all query engines benefit from this rewrite, since some optimizers perform this transformation automatically. Check your specific warehouse's optimizer behavior before rewriting manually.
CTEs vs. subqueries vs. temporary tables
Common Table Expressions (CTEs) are syntactic sugar in many query engines. They do not automatically create materialized intermediate results. If the CTE is referenced multiple times in the outer query, the engine may execute it multiple times rather than computing it once and caching the result.
In Snowflake, CTEs referenced once are typically inlined (treated as subqueries), which is efficient. CTEs referenced more than once may or may not be materialized, depending on the optimizer's cost estimate. In BigQuery, CTEs are generally inlined. In Redshift, the behavior varies.
When a CTE is referenced multiple times in a query and the CTE's computation is expensive, materialize it explicitly as a temporary table. The cost of one additional write is often much less than the cost of computing the CTE twice. This is a warehouse-specific tradeoff: profile before deciding.
Window functions and their performance characteristics
Window functions require sorting the data by the partition and order columns before applying the window computation. For large tables, this sort can dominate query time. A few considerations:
Window functions over an entire table without a PARTITION BY clause force a global sort of the full dataset. Adding a meaningful partition key (such as PARTITION BY user_id) limits the sort to the rows within each partition, which is much cheaper when the partition key has high cardinality.
Multiple window functions that share the same PARTITION BY and ORDER BY clause in a single query can often be computed in a single pass over the sorted data. If you write them as separate subqueries, each may incur an independent sort. Write them together in a single SELECT with matching window definitions to allow the optimizer to compute them together.
Materialized views: where the cost makes sense
A materialized view pre-computes and stores the result of a query. Subsequent queries against the materialized view read the stored result instead of executing the full query. For queries that are expensive to compute and are run frequently against slowly-changing data, materialized views can reduce query time by orders of magnitude.
The cost of a materialized view is freshness: the stored result is stale as soon as the underlying data changes. Warehouses that support incremental materialized view refresh (Snowflake's dynamic tables, BigQuery's incremental refresh) can reduce this staleness, but full refresh materializations may only update hourly or on a schedule you control.
Materialized views are most valuable for: high-cardinality aggregations over large tables that feed interactive dashboards, common join patterns that are referenced across many downstream queries, and complex calculations that take several seconds but are queried hundreds of times per day. They are least valuable for queries that touch recent data (where staleness matters), for one-time exploratory queries, and for queries that vary their filter conditions each time (where the materialized result may not match the specific filter).
When optimizing the query is the wrong answer
Not every slow query is best fixed with query optimization. Some slow queries are slow because the underlying data model is wrong for the access pattern. A star schema optimized for revenue analysis by date and product is not a good data model for a query pattern that primarily accesses data by customer account and last activity date. You can optimize the query to run faster against a model it does not fit, but there is a ceiling on how fast it can go.
When a query remains slow despite reasonable optimization attempts, ask whether the access pattern the query represents should be served by a different physical model: a pre-aggregated summary table, a different sort key, a materialized view at the right grain, or ingestion of the source data into a different destination optimized for that access pattern. Optimizing the SQL is a local fix. Changing the data model is a structural fix that often has much larger returns.
Diagnose slow queries across sources
Nava Labs surfaces query plan details and per-connector timing so you know exactly where time is going in multi-source queries.
Get Early Access