How to Optimize SQL Queries for Data Warehouse Performance: 8 Proven Patterns

Optimizing SQL queries for data warehouse performance requires partitioning tables by date, replacing self-joins with window functions, consolidating aggregations with GROUPING SETS, and pushing filter predicates as early as possible in the execution plan.

Efficient analytical processing in modern cloud data warehouses like Snowflake, Redshift, and BigQuery depends on specific optimization patterns that minimize I/O and maximize parallelism. The DataExpert-io/data-engineer-handbook repository demonstrates concrete implementations showing how to optimize SQL queries for data warehouse performance through strategic schema design and query construction. These techniques reduce scanned rows, eliminate redundant computations, and leverage columnar storage engines effectively.

Design Partitioned Tables for Time-Series Data

Partitioning tables by high-cardinality columns—particularly dates—enables the query engine to perform partition pruning, scanning only relevant data slices rather than full tables.

In intermediate-bootcamp/materials/2-fact-data-modeling/tables/monthly_user_site_hits.sql, the schema defines a composite primary key that includes date_partition:

CREATE TABLE monthly_user_site_hits (
    user_id          BIGINT,
    hit_array        BIGINT[],
    month_start      DATE,
    first_found_date DATE,
    date_partition   DATE,
    PRIMARY KEY (user_id, date_partition, month_start)
);

This structure allows the warehouse to eliminate entire partitions when queries filter on date_partition, dramatically reducing scanned rows. When supported by the warehouse engine, defining primary keys on (user_id, date_partition, month_start) also hints at a clustered layout that optimizes range scans and merge joins.

Replace Self-Joins with Window Functions

Window functions compute cumulative metrics in a single pass through the data, avoiding the costly joins and temporary tables required by self-join approaches.

The file intermediate-bootcamp/materials/4-applying-analytical-patterns/lecture-lab/window_based_analysis.sql demonstrates rolling calculations using SUM(...) OVER (...):

SELECT
    referrer,
    url,
    event_date,
    COUNT(*) AS daily_hits,
    SUM(COUNT(*)) OVER (
        PARTITION BY referrer, url
        ORDER BY event_date
        ROWS BETWEEN 6 PRECEDING AND CURRENT ROW
    ) AS weekly_rolling_count
FROM events_augmented
GROUP BY referrer, url, event_date;

This pattern calculates rolling, monthly, and total counts in one scan, eliminating the need to join the table to itself on date offsets.

Consolidate Multi-Level Aggregations with GROUPING SETS

GROUPING SETS allow a single query to produce multiple aggregation granularities simultaneously, eliminating duplicated table scans for each level.

In intermediate-bootcamp/materials/4-applying-analytical-patterns/lecture-lab/grouping_sets.sql, the following pattern aggregates by OS, device, and browser combinations in one pass:

SELECT
    COALESCE(os_type, '(overall)')   AS os_type,
    COALESCE(device_type, '(overall)') AS device_type,
    COALESCE(browser_type, '(overall)') AS browser_type,
    COUNT(*) AS hits
FROM events_augmented
GROUP BY GROUPING SETS (
    (os_type, device_type, browser_type),
    (os_type),
    (device_type),
    (browser_type)
);

This approach replaces four separate GROUP BY queries with a single execution plan that reads the base table only once.

Push Predicates and Filter Early

Placing WHERE clauses before joins or aggregations allows the optimizer to prune rows early, reducing data movement through the execution pipeline.

While window_based_analysis.sql applies filters like WHERE total_cumulative_sum > 500 after window calculations, optimal performance requires pushing predicates before expensive operations. For example, filtering within the JOIN condition or in a preliminary CTE shrinks the input set before aggregation:

SELECT
    e.user_id,
    d.device_type,
    COUNT(*) AS events_per_device
FROM events e
JOIN devices d ON e.device_id = d.device_id
WHERE e.event_time >= '2023-01-01'   -- filter first
GROUP BY e.user_id, d.device_type;

Filtering on referrer IS NOT NULL before window calculations, as hinted in the repository examples, further minimizes the working dataset.

Optimize Data Types and Avoid Implicit Casts

Selecting appropriate data types preserves columnar compression and prevents the engine from materializing extra representations during query execution.

The monthly_user_site_hits.sql schema uses BIGINT for identifiers and DATE for temporal partitions—types that are natively compressed in most columnar warehouses. Avoiding unnecessary casts in the SELECT list prevents CPU overhead and maintains compression ratios, as casting forces the engine to create additional temporary copies of the data.

Minimize Intermediate Materialization

Common Table Expressions (CTEs) using WITH clauses allow the query optimizer to inline logic when possible, whereas explicit temporary tables add I/O and storage overhead.

The examples in the Data Engineer Handbook favor CTEs over staging tables for intermediate results that are not reused across multiple sessions. This approach enables the engine to eliminate redundant steps through predicate pushdown and avoids the cost of writing intermediate datasets to disk when they are only needed for a single subsequent operation.

Leverage Clustering and Sort Keys

While not explicitly shown with CLUSTER BY syntax in the repository, the primary key definition PRIMARY KEY (user_id, date_partition, month_start) in monthly_user_site_hits.sql suggests a physical layout that warehouses like Snowflake or Redshift can leverage for clustered or sorted storage.

Ordering data on disk by frequently filtered columns enables range scans and faster merge joins. When creating tables, explicitly define sort keys or clustering keys on columns used in JOIN predicates and WHERE clauses to align physical storage with query access patterns.

Limit Result Sets During Exploration

Applying LIMIT clauses during ad-hoc exploration reduces network transfer and can trigger early termination of execution plans before complete aggregation.

While not explicitly illustrated in the analytical pattern files, adding LIMIT 1000 to SELECT statements in window_based_analysis.sql or similar queries accelerates interactive development by preventing the warehouse from processing entire partitions when only sample verification is required.

Summary

  • Partition tables by date columns to enable partition pruning and reduce scanned rows.
  • Use window functions instead of self-joins for cumulative and rolling metrics to avoid expensive join operations.
  • Implement GROUPING SETS to generate multiple aggregation levels in a single table scan.
  • Push filter predicates as early as possible, preferably before joins and window functions.
  • Select native data types like BIGINT and DATE to maintain columnar compression and avoid casting overhead.
  • Prefer CTEs over temporary tables for single-use intermediate logic to minimize I/O.
  • Define clustering keys on frequently filtered columns to optimize range scans and join performance.
  • Apply LIMIT clauses during development to reduce execution time for exploratory queries.

Frequently Asked Questions

What makes window functions more efficient than self-joins for time-series analysis?

Window functions scan the table once and maintain a rolling window in memory, whereas self-joins require the database to read the table multiple times and match rows based on complex join conditions. As implemented in window_based_analysis.sql, a single SUM(...) OVER (...) clause replaces multiple self-joins that would otherwise create intermediate Cartesian products and temporary tables.

When should I use GROUPING SETS instead of UNION ALL?

Use GROUPING SETS when you need multiple aggregation granularities from the same base data, such as totals by day, month, and year. This approach, demonstrated in grouping_sets.sql, reads the source table only once, whereas UNION ALL would execute separate scans for each aggregation level, multiplying I/O costs.

How does predicate pushdown improve data warehouse performance?

Predicate pushdown moves filter conditions as close to the data source as possible, often eliminating rows before expensive operations like joins, sorts, or window functions. In columnar warehouses, this reduces the amount of compressed data read from storage and minimizes network traffic between compute nodes, as shown in the early-filtering join patterns from the repository.

Should I use CTEs or temporary tables for complex queries?

Use CTEs for intermediate results that are consumed immediately within the same query and not reused across sessions. CTEs allow the optimizer to inline the logic and eliminate unnecessary steps. Temporary tables are preferable only when the intermediate result is large, complex, and referenced multiple times in separate queries, as materialization adds write and read overhead.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →