skip to content

How does Amazon Redshift execute a join between a local table and a Spectrum external table?

level: seniorimportance: should knowfreq 50%

answer

  1. the work is split across two tiers
  2. scan and filter out there, join in here
  3. external rows arrive with no home slice
  4. broadcast versus redistribute, and who decides
  5. the planner has no idea how big the S3 table is

basics

~20 s

Spectrum's scan layer reads, projects, filters and often partially aggregates the S3 data, then streams the surviving rows to the cluster's compute nodes. The join itself always runs on the cluster, and because external tables have no distribution key those rows must be broadcast or redistributed first.

solid answer

~50 s

The work splits across two tiers. The Spectrum fleet handles the S3 scan and pushes down projection, comparison and pattern predicates, `GROUP BY`, `DISTINCT`, `LIMIT` and the common aggregates. **Joins are never pushed down** — they execute on your compute nodes, along with window functions and final aggregation. Since an external table carries no `DISTKEY` or `SORTKEY`, the rows arriving from S3 have no useful placement, so the plan shows `DS_BCAST_INNER` or a redistribution to line them up with the local table. `EXPLAIN` makes this visible as `S3 Seq Scan`, `S3 HashAggregate` and `S3 Query Scan` nodes beneath an `XN Hash Join`. Two things keep it fast: shrink the external side in S3 first (partition predicate, tight projection, pre-aggregation) so few rows cross, and give the planner a row-count hint via `SET TABLE PROPERTIES ('numRows'='...')` so it does not broadcast the wrong side.

code

text · 7 lines
text
XN Hash Join DS_BCAST_INNER
  ->  XN S3 Query Scan events
        ->  S3 HashAggregate
              ->  S3 Seq Scan spectrum.events
                    location:"s3://my-lake/events" format:PARQUET
  ->  XN Hash
        ->  XN Seq Scan on dim_user

go deeper

for a junior

Know that the S3 side is read and filtered by a separate Spectrum layer, and that the join itself happens on the Redshift cluster after those rows arrive.

for a middle

Explain what gets pushed down — projection, predicates, GROUP BY, DISTINCT, LIMIT, the basic aggregates — and why joins do not, and what that means for how much data crosses the network.

for a senior

Demonstrate reading an EXPLAIN for DS_BCAST_INNER versus redistribution, setting numRows on the external table, and restructuring the query to pre-aggregate in S3 before the join.

for a principal

Own the boundary decision: which joins are allowed to reach S3 at all, when a materialized view or a local hot table replaces the pattern, and how that shows up in scan spend.

## Two tiers of execution A query mixing local and external tables runs in two places. **The Spectrum layer** is an AWS-managed scan fleet that reads your S3 objects. Redshift pushes work into it aggressively: column projection, comparison predicates (`=`, `<`, `>`, `BETWEEN`), pattern predicates like `LIKE`, `GROUP BY` clauses, `DISTINCT`, `LIMIT`, and aggregate functions such as `COUNT`, `SUM`, `AVG`, `MIN` and `MAX`. That means the fleet can often return a pre-aggregated result rather than raw rows. **The cluster** — your leader and compute nodes — does everything else. Crucially that includes **all joins**, plus window functions and any aggregation that must combine partial results. ## Why the join costs what it costs A local Redshift table has a distribution style, so its rows sit on a known slice. An external table has none: it is a stream of rows arriving from an outside fleet with no relationship to your slices. To join it to a local table, Redshift must first co-locate matching rows, and it has only two ways: - **Broadcast** (`DS_BCAST_INNER`): send a full copy of the smaller side to every node. Cheap when that side is genuinely small — a dimension of thousands of rows. - **Redistribute** (`DS_DIST_BOTH` / `DS_DIST_INNER`): hash both sides on the join key and shuffle them across the network. Necessary when both sides are large, and expensive. A plan looks roughly like this: ```text XN Hash Join DS_BCAST_INNER -> XN S3 Query Scan events -> S3 HashAggregate -> S3 Seq Scan spectrum.events location:"s3://my-lake/events" format:PARQUET -> XN Hash -> XN Seq Scan on dim_user ``` Read it bottom-up on the external branch: scan S3, aggregate in the Spectrum layer, hand the result to the cluster, join. ## The statistics problem The planner has no statistics for an external table. Without them it may guess badly and broadcast a 4-billion-row external side across every node — the classic Spectrum blow-up. Redshift lets you tell it: ```sql ALTER TABLE spectrum.events SET TABLE PROPERTIES ('numRows'='4200000000'); ``` Setting a realistic `numRows` (and refreshing it as the table grows) is one of the highest-leverage things you can do for mixed joins, because it drives join order and distribution choice. ## Shrink the external side first The governing principle: **do as much as possible in the Spectrum layer, and let as little as possible cross to the cluster.** Concretely: - Always put a partition predicate on the external table so fewer files are opened. - Project only the columns the join and the output need. - Pre-aggregate in a subquery over the external table so the join operates on a summary, not on raw events. If the query ultimately produces daily totals per user, aggregate to that grain in the external scan and join the small result to the dimension. - Keep the local dimension small enough to broadcast, or filter it before the join. The anti-pattern is `SELECT * FROM spectrum.big_table b JOIN local_dim d ON ...` with the only selective filter sitting on `d`. The filter lives on the wrong side: Spectrum cannot use it, so the entire external table is read and streamed. ## When to stop joining externally If the same external-to-local join runs repeatedly, stop paying for it repeatedly. Options, in escalating order: 1. A **materialized view over the external table**, refreshed on a schedule, so queries hit local storage. 2. `COPY` the hot partitions into a **local table** with a `DISTKEY` matching the join column and a `SORTKEY` on the filter column — now the join is co-located and there is no S3 scan at all. 3. Keep only cold history external, joined rarely, in a `UNION ALL` view with the local hot table. ## What to measure `SVL_S3QUERY_SUMMARY` shows `s3_scanned_bytes` versus `s3query_returned_bytes` — the ratio tells you how much filtering the Spectrum layer actually did. If returned bytes are large, the external side was not shrunk and the network transfer plus redistribution is your real cost. `EXPLAIN` tells you which side is being broadcast; if it is the external one, your `numRows` property is probably missing or wrong.

  • Which parts of a query does Redshift push down into the Spectrum layer, and which never go there?
    Pushed down: column projection, comparison and pattern predicates, `GROUP BY`, `DISTINCT`, `LIMIT`, and aggregates like `COUNT`, `SUM`, `AVG`, `MIN` and `MAX`. Never pushed down: joins, window functions, and any final combination of partial aggregates. That asymmetry is why you want the external side reduced to a summary before it reaches the cluster.
  • Why does setting numRows on an external table matter so much?
    External tables have no statistics, so the planner guesses cardinality. A bad guess leads it to broadcast a huge external side to every node, which turns a fast query into an hours-long one. `ALTER TABLE ... SET TABLE PROPERTIES ('numRows'='...')` gives the optimiser a realistic figure so it picks join order and distribution sensibly.
  • The only selective filter in a mixed join sits on the local dimension. Why is that a problem?
    Spectrum can only push down predicates written against the external table. A filter on the local side is applied after the external rows have already been read from S3 and streamed to the cluster, so you pay full scan cost. Either express an equivalent predicate on the external table's own columns, or resolve the dimension first and inline the resulting key list.

saying these in an interview costs you the question

  • Believing joins are pushed down to the Spectrum layer
  • Expecting a DISTKEY on an external table to co-locate the join
  • Assuming a filter on the local side reduces S3 bytes scanned
  • Ignoring numRows and letting the planner broadcast the external side
  • Thinking a slow mixed join is always an S3 throughput problem

context