+1 (726) 227-3497

Reading Redshift Query Plans: EXPLAIN, DS_BCAST_INNER and the Rewrites That Fix Slow SQL

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:

LabelMeaningVerdict
DS_DIST_NONENo movement — both sides already colocatedBest
DS_DIST_ALL_NONEInner table is DISTSTYLE ALL, already on every nodeGood
DS_DIST_INNERInner table redistributed on the join keyAcceptable if inner is small
DS_DIST_OUTER / DS_DIST_BOTHOne or both sides reshuffled across the clusterInvestigate
DS_BCAST_INNEREntire inner result copied to every nodeRed 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:

  1. Fix the statistics. Run ANALYZE on 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.
  2. Make the small side DISTSTYLE ALL. For dimensions under a few million rows that are joined constantly, ALL turns every join into DS_DIST_ALL_NONE at the cost of storing a copy per node.
  3. Colocate the big-to-big join. Two large tables joined on the same key should share a DISTKEY on that key, giving DS_DIST_NONE. You only get one distribution key per table, so spend it on the join that dominates your workload.
  4. 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 the OVER().
  • Replace SELECT DISTINCT over wide rows with GROUP BY on 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

  1. Find the regression in SYS_QUERY_HISTORY — compare this week's p95 elapsed_time per query_text fingerprint against last week's.
  2. Pull SYS_QUERY_DETAIL for a slow execution. Look for is_diskbased, is_rrscan = 'f', and slice skew.
  3. EXPLAIN the statement. Find the first DS_BCAST_INNER, DS_DIST_BOTH or Nested Loop from the bottom up — the deepest bad step poisons everything above it.
  4. ANALYZE the tables involved and re-EXPLAIN before changing any DDL. Half of all bad plans are stale statistics.
  5. Apply one change. Re-run. Record the before/after elapsed_time and output_bytes in 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.