Redshift's SYS views tell you which query was slow. They do not tell you why. For that you need the query plan — and on an MPP engine the plan contains one category of information that single-node databases do not have: how much data has to move between compute nodes before the join can run. Most "mysteriously slow" Redshift queries are not slow because of CPU or disk. They are slow because the optimizer decided to broadcast a 400-million-row intermediate result to every node.
This tutorial covers reading EXPLAIN output, mapping the estimated plan to what actually happened with SVL_QUERY_SUMMARY and the SYS_QUERY_DETAIL views, and the handful of rewrites that fix the common cases.
The anatomy of a Redshift plan
EXPLAIN gives you the optimizer's estimate, bottom-up:
EXPLAIN
SELECT c.segment, SUM(o.net_amount) AS revenue
FROM sales.orders o
JOIN sales.customers c ON c.customer_id = o.customer_id
WHERE o.order_date >= '2026-01-01'
GROUP BY 1;
XN HashAggregate (cost=18244452.01..18244452.07 rows=5 width=28)
-> XN Hash Join DS_BCAST_INNER (cost=48112.00..18244001.11 rows=90178 width=28)
Hash Cond: ("outer".customer_id = "inner".customer_id)
-> XN Seq Scan on orders o (cost=0.00..902105.00 rows=88104210 width=12)
Filter: (order_date >= '2026-01-01'::date)
-> XN Hash (cost=48112.00..48112.00 rows=4811200 width=24)
-> XN Seq Scan on customers c (cost=0.00..48112.00 rows=4811200 width=24)
Three things matter far more than the cost numbers.
1. The data-movement label on every join. Redshift annotates joins with how rows are redistributed:
| Label | Meaning | Verdict |
|---|---|---|
DS_DIST_NONE | No movement — both sides already colocated | Best |
DS_DIST_ALL_NONE | Inner table is DISTSTYLE ALL, already on every node | Good |
DS_DIST_INNER | Inner table redistributed on the join key | Acceptable if inner is small |
DS_DIST_OUTER / DS_DIST_BOTH | One or both sides reshuffled across the cluster | Investigate |
DS_BCAST_INNER | Entire inner result copied to every node | Red flag at scale |
A broadcast of 50,000 rows is free. A broadcast of 50 million rows on a 10-node cluster means moving half a billion rows over the network before any join work starts.
2. The join algorithm. Hash Join and Merge Join are fine. Nested Loop almost never is — it usually means you wrote a cross join by accident, or your join predicate is an inequality or a function the optimizer cannot hash.
3. The row estimates. rows=5 for the aggregate above is a guess. If the estimates are wildly wrong, every downstream decision — broadcast vs. redistribute, join order, hash table sizing — is built on sand. Bad estimates almost always mean stale statistics.
Estimated plan vs. what actually ran
EXPLAIN lies by omission. Confirm against execution. On provisioned and Serverless, the SYS_ views are the current interface:
-- Step-level detail for one query
SELECT step_name, table_name,
input_bytes / 1048576.0 AS input_mb,
output_bytes / 1048576.0 AS output_mb,
input_rows, output_rows,
is_rrscan, is_diskbased,
duration / 1000000.0 AS seconds
FROM SYS_QUERY_DETAIL
WHERE query_id = 1234567
ORDER BY step_name;
Two columns are worth alerting on:
is_diskbased = 't'— a hash, sort or aggregate spilled to disk because it did not fit in memory. This is the single most common cause of a query that is 10x slower than its neighbours.is_rrscan = 'f'on a large scan — no range-restricted scan, meaning Redshift read every block instead of skipping via zone maps. Your filter is not aligned with the sort key, or the predicate is wrapped in a function.
Skew is the other thing the plan will not show you. Check whether one slice did all the work:
SELECT query_id, step_name,
MAX(output_rows) AS max_slice_rows,
AVG(output_rows) AS avg_slice_rows,
MAX(output_rows) / NULLIF(AVG(output_rows), 0) AS skew
FROM SYS_QUERY_DETAIL
WHERE query_id = 1234567
GROUP BY 1, 2
HAVING MAX(output_rows) / NULLIF(AVG(output_rows), 0) > 3
ORDER BY skew DESC;
A skew ratio above ~3 on a join or aggregate step means your distribution key (or the join key) has a dominant value — frequently a NULL, a -1 sentinel, or 'unknown' from an unmatched dimension lookup.
Fix 1: stop the broadcast
If a large inner side is being broadcast, you have four levers, in order of preference:
- Fix the statistics. Run
ANALYZEon both tables and re-EXPLAIN. A table that has grown 50x since its last analyze will be estimated at its old size, and the optimizer will cheerfully broadcast it. - Make the small side
DISTSTYLE ALL. For dimensions under a few million rows that are joined constantly,ALLturns every join intoDS_DIST_ALL_NONEat the cost of storing a copy per node. - Colocate the big-to-big join. Two large tables joined on the same key should share a
DISTKEYon that key, givingDS_DIST_NONE. You only get one distribution key per table, so spend it on the join that dominates your workload. - Shrink the inner side before the join. Push filters and aggregation into a CTE or subquery so the thing being redistributed is small:
WITH recent_customers AS (
SELECT customer_id, segment
FROM sales.customers
WHERE status = 'active'
)
SELECT rc.segment, SUM(o.net_amount)
FROM sales.orders o
JOIN recent_customers rc ON rc.customer_id = o.customer_id
WHERE o.order_date >= '2026-01-01'
GROUP BY 1;
With Automatic Table Optimization on, check what Redshift already wants to change before you hand-tune anything:
SELECT type, database, table_name, ddl, auto_eligible
FROM SVV_ALTER_TABLE_RECOMMENDATIONS;
Fix 2: make filters zone-map friendly
Redshift stores min/max zone maps per 1 MB block. A predicate can only skip blocks if it is a direct comparison against the sort key column:
-- Bad: function on the column defeats the zone map
WHERE DATE_TRUNC('month', order_date) = '2026-03-01'
-- Good: a range the zone map can use
WHERE order_date >= '2026-03-01' AND order_date < '2026-04-01'
The same applies to implicit casts (WHERE order_date::text LIKE '2026-03%') and to joining a date column to a string. Check the result in SYS_QUERY_DETAIL: is_rrscan should flip to t and the scan's input_rows should drop sharply.
Fix 3: kill the nested loop
XN Nested Loop DS_BCAST_INNER (cost=0.00..441520983.00 rows=81249030 width=44)
A nested loop with no join condition in the plan means the predicate was lost — a missing ON clause, or a join written in the WHERE clause with an OR that the optimizer cannot decompose. Range joins (BETWEEN on effective dates, IP-range lookups) also force nested loops. The standard fix is to add an equality component the optimizer can hash on — a bucket key:
-- Range join made hashable by bucketing on month
SELECT f.*, r.rate
FROM facts f
JOIN rates r
ON r.month_key = DATE_TRUNC('month', f.event_at) -- equality: hashable
AND f.event_at >= r.valid_from
AND f.event_at < r.valid_to; -- residual filter
Set query_group/abort rules so these never run unbounded in production. A WLM query-monitoring rule with nested_loop_join_row_count > 1000000 → abort catches accidental cross joins before they eat the cluster.
Fix 4: stop the spill
is_diskbased = 't' on a hash or sort step means the step did not fit in the slice's memory allocation. Options, in order:
- Reduce what you are sorting. A
SELECT *carried through a window function materialises every column; project only what you need before theOVER(). - Replace
SELECT DISTINCTover wide rows withGROUP BYon the key columns. - Give the query more memory for a single run —
SET wlm_query_slot_count TO 3;on a manual WLM queue, or route it to a queue/priority with a larger share under automatic WLM. - If a big step spills every time by design (an annual backfill, a full re-aggregation), pre-materialise it: an incremental materialized view or a staged intermediate table is cheaper than repeatedly spilling.
A repeatable triage loop
- Find the regression in
SYS_QUERY_HISTORY— compare this week's p95elapsed_timeperquery_textfingerprint against last week's. - Pull
SYS_QUERY_DETAILfor a slow execution. Look foris_diskbased,is_rrscan = 'f', and slice skew. EXPLAINthe statement. Find the firstDS_BCAST_INNER,DS_DIST_BOTHorNested Loopfrom the bottom up — the deepest bad step poisons everything above it.ANALYZEthe tables involved and re-EXPLAINbefore changing any DDL. Half of all bad plans are stale statistics.- Apply one change. Re-run. Record the before/after
elapsed_timeandoutput_bytesin the ticket, so the next person knows the plan was deliberate.
Resist the temptation to fix five things at once. On an MPP engine a single distribution-key change can move a query from 40 minutes to 40 seconds, and if you changed the sort key, the DISTKEY and three predicates in the same deploy, you will never know which one did it — or which one you can safely revert.
When to bring in help
Plan-level tuning pays off fastest when someone has seen the same failure modes across many clusters. If you have a dashboard that times out, a nightly job whose runtime has quietly tripled, or a cluster you keep resizing because queries "feel slow", our Redshift Performance Optimization practice does exactly this work — and the findings usually pay for themselves in right-sized compute. Get in touch with a query ID and we will start from the plan.