What does DS_BCAST_INNER in a Redshift EXPLAIN plan tell you, and how do you remove it?
answer
- The plan tells you what had to move
- One side got copied to every node
- Cost scales with node count, not just size
- Co-locating on the join column removes it
- Small dimensions can be replicated at load time instead
basics
~20 sDS_BCAST_INNER means Redshift is copying the entire inner table to every compute node before the join, because the two sides are not co-located. You remove it by giving both tables the same DISTKEY as the join column, or by making the small side DISTSTYLE ALL.
solid answer
~50 sRedshift annotates each join in `EXPLAIN` with how it had to move data. `DS_BCAST_INNER` says the whole inner table was broadcast to every node at query time — the plan's most expensive distribution label. `DS_DIST_BOTH` means both sides were reshuffled on the join key, also expensive. `DS_DIST_INNER` reshuffles only the inner side. `DS_DIST_NONE` and `DS_DIST_ALL_NONE` are what you want: no movement at all, because the rows were already on the right slice. The cure depends on the tables. If both sides are large, distribute both on the join column so matching rows land on the same slice. If the inner side is a small dimension, `DISTSTYLE ALL` makes it locally available everywhere and the plan flips to `DS_DIST_ALL_NONE`. Also check that statistics are fresh — a stale row estimate can make Redshift broadcast a table that is no longer small.
code
text · 6 linesXN HashAggregate (cost=...)
-> XN Hash Join DS_BCAST_INNER (cost=...)
Hash Cond: ("outer".customer_id = "inner".customer_id)
-> XN Seq Scan on fact_orders o (cost=...)
-> XN Hash (cost=...)
-> XN Seq Scan on dim_customer c (cost=...)go deeper
Know that a Redshift plan can show data being moved between nodes before a join, and that DS_DIST_NONE means no movement was needed.
Explain each DS_ label, why broadcast cost multiplies by node count, and the two structural fixes: matching DISTKEYs on both sides, or DISTSTYLE ALL on the small side.
Diagnose from a real plan under time pressure — distinguish a genuine distribution mismatch from a stale-statistics misestimate, and know when redistribution is simply acceptable.
Reason about the whole join graph: which one join each fact table spends its single DISTKEY on, which dimensions justify replication, and what that costs in storage and load time across the schema.
## What the label is Running `EXPLAIN` on a Redshift query prints a plan whose join nodes carry a distribution attribute — a `DS_` prefixed label describing what the engine must do to get the two sides of the join onto the same slice before it can match rows. Reading those labels is the fastest way to find out why an MPP join is slow, and they show up in almost every real Redshift tuning conversation. The labels you will meet: - **`DS_DIST_NONE`** — no movement. Both tables are distributed on the join column, so matching rows already share a slice. This is the goal. - **`DS_DIST_ALL_NONE`** — no movement, because the inner table is `DISTSTYLE ALL` and therefore already present on every node. Equally good. - **`DS_DIST_INNER`** — the inner table is redistributed across slices on the join key. The outer side stays put. Moderate cost, proportional to the inner table's size. - **`DS_DIST_OUTER`** — the outer table is redistributed instead. - **`DS_BCAST_INNER`** — a full copy of the inner table is sent to every compute node. Cost scales with inner table size *times* node count. - **`DS_DIST_BOTH`** — both tables are redistributed on the join key. Usually the worst case on large inputs, since everything crosses the network. ## Why broadcasting happens A join can only compare two rows that sit on the same slice. When the distribution keys don't line up with the join predicate, Redshift has two ways to fix that at runtime: redistribute (rehash both sides on the join column and ship rows to their new owners) or broadcast (replicate one side everywhere). The optimizer picks broadcast when it believes one side is small enough that copying it wholesale is cheaper than reshuffling both. That belief comes from statistics. If `ANALYZE` has not run since a table grew from 10,000 rows to 100 million, the planner may still think it is tiny and cheerfully broadcast it — which is why a query that was fine last quarter can suddenly take an order of magnitude longer with no change to its text. A `DS_BCAST_INNER` on a table you know is large is a strong hint to check `stats_off` in `SVV_TABLE_INFO` and re-run `ANALYZE`. ## Fixing it structurally **Large joined to large.** Distribute both tables on the join column. If `fact_orders` and `fact_order_items` both join on `order_id`, `DISTKEY (order_id)` on both turns the join into `DS_DIST_NONE`. The constraint is that each table has only one DISTKEY, so this works for the one join you care most about. **Large joined to small.** Set the small side to `DISTSTYLE ALL`. Every node already has the whole dimension, so the plan becomes `DS_DIST_ALL_NONE` regardless of how the fact table is distributed — and unlike a broadcast, the copying happens once at load time rather than on every query execution. This is the single most common fix in a star schema. **Neither.** Sometimes the join is rare and the redistribution is acceptable. `DS_DIST_INNER` on a modest table once an hour is not worth reshaping a schema for. Tune what the workload actually runs. ## What broadcast actually costs Broadcast cost is not just the bytes moved; it is bytes multiplied by node count, and the copies land in each node's memory. On a wide table with large `VARCHAR` columns this can push the join into spilling to disk, which turns a network problem into an I/O problem. A useful mitigation independent of distribution is to reduce what gets broadcast — project only the columns you need and push filters below the join so the inner side is genuinely small before it is copied. ## Reading it in practice ```sql EXPLAIN SELECT c.segment, SUM(o.amount) FROM fact_orders o JOIN dim_customer c ON c.customer_id = o.customer_id WHERE o.order_ts >= '2026-01-01' GROUP BY 1; ``` If the join line reads `XN Hash Join DS_BCAST_INNER`, `dim_customer` is being copied to every node on each run. Setting `dim_customer` to `DISTSTYLE ALL`, or giving both tables `DISTKEY (customer_id)`, changes that line to `DS_DIST_ALL_NONE` or `DS_DIST_NONE` respectively. One caution: `EXPLAIN` shows the *estimated* plan. To see what actually happened — including how many bytes really crossed the network and whether any step spilled — look at the execution-time system views such as `SVL_QUERY_REPORT` and `SVL_QUERY_SUMMARY` for the query id. Estimated and actual can diverge exactly when statistics are stale, which is the case you are usually hunting.
- How does DS_BCAST_INNER differ from DS_DIST_BOTH in cost?Broadcast copies the entire inner table to every node, so cost is inner size times node count, and each copy occupies node memory. DS_DIST_BOTH rehashes both sides on the join key and each row moves once, so total bytes moved is roughly the sum of both tables. Broadcast wins when the inner side is genuinely small; redistribution wins when both are large.
- Why can DISTSTYLE ALL beat a broadcast that moves the same rows?Both put the whole dimension on every node, but ALL does the copying once at load time and stores it, while a broadcast repeats the copy on every query execution and holds it in query memory. If the table is joined constantly, paying the duplication once in storage is far cheaper than paying it per query.
- You see DS_BCAST_INNER on a table you know has 400 million rows. What do you check first?Statistics. The optimizer only broadcasts what it thinks is small, so check stats_off in SVV_TABLE_INFO for that table and run ANALYZE. A table that grew dramatically since its last analyze is the classic cause of a plan that suddenly starts broadcasting something huge.
- Does an aggregation without a join ever show these distribution labels?Yes — a GROUP BY whose keys don't match the table's DISTKEY also needs rows brought together, and the plan will show a redistribution step feeding the aggregate. Distributing a table on a column you group by heavily lets Redshift aggregate locally per slice first.
saying these in an interview costs you the question
- Reading DS_BCAST_INNER as a good sign because 'broadcast joins are fast'
- Thinking a sort key can prevent redistribution at join time
- Assuming broadcast cost depends only on table size, not node count
- Ignoring stale statistics as the reason a large table gets broadcast
- Believing you can fix any join by adding a second DISTKEY column