Query performance debugging in column-oriented analytics databases follows a different diagnostic path than tuning queries in PostgreSQL or MySQL. The physical storage model, the query optimizer, and the bottleneck locations are different. If you are using intuitions from OLTP tuning and applying them to Redshift, BigQuery, Snowflake, or DuckDB, some of those intuitions will be wrong.
This article covers a diagnostic sequence for slow queries in columnar analytics databases, organized by the failure mode that explains the slowness.
Start with EXPLAIN, but know its limits
The first step in diagnosing a slow query is always reading the query execution plan. EXPLAIN or EXPLAIN ANALYZE (or its equivalent in your specific warehouse) tells you what the optimizer decided to do. For slow queries, you are looking for a few specific signals.
Row count estimates that are dramatically wrong are a leading indicator of plan quality problems. Most columnar optimizers use statistics to estimate how many rows each node in the plan will process. When estimates are off by an order of magnitude or more (the optimizer expects 1,000 rows, actually gets 10 million), the plan choices built on those estimates can be badly suboptimal: wrong join order, wrong join algorithm, wrong amount of memory allocated for sorting.
The problem is that EXPLAIN tells you the plan but not whether that plan is optimal. A plan where the optimizer chose a merge join over a hash join might be correct or might be a consequence of outdated column statistics. You need to understand what the right plan should look like before you can evaluate whether the actual plan is sensible.
Partition pruning: are you scanning too much data?
For large partitioned tables, the first thing to verify is whether partition pruning is working. Partition pruning means the query engine skips partitions that cannot contain rows matching the query's filter conditions. If your orders table is partitioned by month and your query filters on WHERE order_date >= '2025-10-01' AND order_date < '2025-11-01', the engine should scan only the October 2025 partition, not the entire table.
When partition pruning fails, the query scans the full table regardless of the date filter. This is one of the most common sources of "this query used to be fast and now it is slow" complaints: the table grew from 3 months of data to 24 months, partition pruning broke silently (often due to a type mismatch or implicit cast on the filter expression), and now what was a 1-partition scan is a 24-partition scan.
To verify partition pruning, check the query plan for a "partition filter" or "dynamic partition pruning" node, or check the metadata in EXPLAIN ANALYZE for how many partitions were scanned. If the number of partitions scanned does not match your expectation given the filter, pruning is broken. The usual fix is ensuring the filter expression uses the exact type of the partition column without implicit casts. order_date::DATE on a column stored as TIMESTAMP may defeat pruning in some systems.
Join ordering and broadcast vs. hash join selection
Join performance in distributed columnar systems depends heavily on the join strategy the optimizer picks. The two most important strategies are broadcast joins and hash (shuffle) joins.
A broadcast join sends a copy of the smaller table to every worker node, allowing each worker to join its partition of the larger table against the full smaller table locally. This is fast when the smaller table is genuinely small. A hash join partitions both tables by the join key and sends matching partitions to the same worker. This works for large-to-large joins but requires a full shuffle of the data across the cluster.
The optimizer decides which strategy to use based on table size statistics. When those statistics are stale, the optimizer may choose a broadcast join for a table that is now too large to broadcast efficiently, or choose a hash join for a table that could be broadcast. Both choices degrade performance.
The diagnostic: find the join nodes in the plan and check whether the optimizer's size estimates match reality. In Redshift, ANALYZE refreshes table statistics. In Snowflake, the warehouse maintains statistics automatically but occasionally they lag for very recently loaded data. If you suspect stale statistics, run the statistics refresh command and re-run the query.
Predicate pushdown in federated and multi-layer queries
For queries that span multiple layers (a federated query that reads from an operational database and joins against a warehouse table, or a query through a logical layer like a dbt model), predicate pushdown is critical and often fails silently.
Predicate pushdown means that a filter in the outer query (the one your analytics tool sends) is pushed down through the query layers to execute as close to the storage as possible. If your dashboard query filters on WHERE region = 'US-West', that filter should be pushed into the SQL that runs against the underlying table, limiting the rows returned at the source. If pushdown fails, the full unfiltered table is scanned and the filter applies only after all rows have been transferred to the query execution layer.
Pushdown failures happen when intermediate view definitions or transformation layers use constructs that block optimizer transparency: DISTINCT in a subquery, LIMIT without ORDER BY, or window functions without a matching partition key. When you add a filter to a query against a view that uses any of these constructs, the filter cannot be pushed through the view into the underlying table scan.
The fix is to check whether your intermediate transformation layers (dbt models, views, materialized views) are written in a way that enables filter pushdown. Replacing blocking patterns with alternatives, or materializing the intermediate result as a table rather than a view, can restore pushdown and dramatically reduce the data scanned.
Spill to disk: when queries run out of memory
In-memory query execution is the normal fast path. When a query requires more memory than the warehouse has allocated for it (typically during large sort operations or hash joins on large tables), it spills intermediate data to disk. Spill to disk is often 10x to 100x slower than in-memory execution.
Queries that run fine on small datasets and become extremely slow as data grows are often hitting a spill boundary. The query ran in memory at 50 million rows and spills at 200 million rows. The runtime difference is not linear; it is often a step function where crossing the memory boundary produces a sudden dramatic slowdown.
Most analytics warehouses expose spill metrics in their query execution logs or the EXPLAIN ANALYZE output. Snowflake's query profile UI shows bytes spilled to local and remote storage. Redshift's STL_QUERY_METRICS view shows per-step spill information. If your slow query is spilling, the fix is either to increase the warehouse size (more memory per node) or to refactor the query to reduce peak memory usage: filter before joining, aggregate before joining, break a large query into smaller steps that each complete in memory.
Skew: when some workers do all the work
Data skew is one of the harder-to-diagnose sources of query slowness in distributed systems. A well-distributed query has all worker nodes finishing their assigned work at roughly the same time. A skewed query has one or a few worker nodes doing much more work than the rest, while all other nodes finish and wait.
Join key skew is the most common case. If your query joins on customer_id and a small number of customer IDs account for a large fraction of the rows (a few large accounts that produce orders of magnitude more activity than typical customers), all the rows for those customer IDs get routed to the same worker node during a hash join. That node does 30 percent of the work while all other nodes do 2-3 percent each.
The indicator is in the per-node or per-step timing: if one step shows a very long elapsed time and the per-node breakdown shows one node taking 10x longer than average, skew is the cause. The fix depends on the warehouse but generally involves either a hash randomization technique (salting the join key for the skewed values to spread them across nodes) or moving to a broadcast join for the table with skewed keys.
A diagnostic sequence to follow
Given a slow query, the sequence that catches the most common issues in order:
- Run
EXPLAIN ANALYZEand check the estimated vs. actual row counts at each node. Stale statistics = refresh and re-run. - Check partition scan count against your expectation for the filter. Unexpected full scan = check for implicit casts on the partition column.
- Find the largest scans and joins in the plan. Are size estimates reasonable? Are join strategies appropriate for the estimated table sizes?
- Check for spill in execution logs. Spill present = reduce peak memory usage or increase warehouse tier.
- Check per-node timing for skew. One node taking 10x longer than others = join key skew investigation.
This sequence covers a large fraction of slow query root causes in columnar analytics databases. The specifics of where to find each signal differ by warehouse (BigQuery's execution details, Snowflake's query profile, Redshift's explain analyzer), but the underlying failure modes are consistent across systems.
Query across sources with full plan visibility
Nava Labs exposes query plans and per-connector execution stats so you can see where time goes across federated queries.
Get Early Access