skip to content

In an Amazon Redshift provisioned cluster, what does the leader node do that compute nodes do not?

level: juniorimportance: must knowfreq 80%

answer

  1. One node talks to clients
  2. Planning happens in exactly one place
  3. User data lives only on compute nodes
  4. Final merge is single-threaded on the leader

basics

~20 s

The leader node is the only node clients connect to: it parses SQL, builds and compiles the execution plan, ships code to the compute nodes, and merges their partial results. Compute nodes store the data and execute plan segments in parallel.

solid answer

~50 s

A provisioned Redshift cluster is one leader node plus one or more compute nodes. The leader owns the client endpoint and everything single-threaded about a query: parsing, catalog lookup, optimization, compiling plan segments into executable code, distributing that code to the compute nodes, and performing the final merge of the partial results they send back. It stores **no user table data**. Compute nodes are divided into slices; each slice holds a portion of every table's rows and runs the plan segments over its own data in parallel, exchanging rows with other nodes when a join or aggregation needs it. Queries that touch only catalog tables run entirely on the leader. The practical consequence is that anything which funnels many rows through the leader — an unaggregated `SELECT *`, a final `ORDER BY` over a huge result — serializes on one machine and is where naive Redshift queries fall over.

code

sql · 7 lines
sql
-- Every row funnels through the single leader node
SELECT * FROM events ORDER BY event_ts;

-- Compute nodes aggregate in parallel; the leader merges a few rows
SELECT event_type, COUNT(*)
FROM events
GROUP BY event_type;

go deeper

for a junior

Be able to draw the box diagram: one leader node the client connects to, several compute nodes holding the data. Know that the leader stores no user table data.

for a middle

Explain what the leader actually produces — a parsed, optimized, compiled plan — and how compiled-code caching creates a slow first run. Describe how compute nodes execute segments over their own slices.

for a senior

Show that you diagnose from the architecture: identify plan shapes whose cost lands on the single leader (unaggregated result sets, final sorts, connection counts) and prescribe the fix, such as aggregating in-cluster or unloading to S3.

for a principal

Frame the leader as the cluster's serialization point and a shared resource across every workload on it, and reason about when that argues for splitting workloads across separate compute rather than adding nodes to one cluster.

## What a cluster is made of A provisioned Amazon Redshift cluster consists of exactly one **leader node** and one or more **compute nodes** connected by a fast private network. Clients — JDBC/ODBC drivers, BI tools, `psql` — connect only to the leader's endpoint; compute nodes are not reachable from outside. In a multi-node provisioned cluster you are billed for compute node hours; the leader node hours are not charged. ## What the leader node does The leader node performs all the coordination work of a query: - **Parse and validate.** It parses the SQL and resolves table and column names against the system catalog it maintains. - **Optimize.** It produces a distributed execution plan: join order, join method, and — crucially for an MPP engine — whether each join input must be broadcast to every node or redistributed on the join key. - **Compile.** Redshift does not interpret plans row by row; it generates and compiles native code for the plan's segments. Compiled segments are cached, so the *first* execution of a brand-new query shape can be noticeably slower than the second. A large scale-out compilation service backs this, but the effect still shows up in cold benchmarks. - **Distribute.** It sends the compiled segments to every compute node, along with any small pieces of data the plan needs. - **Merge and return.** It receives partial results from the nodes, performs the final aggregation, merge-sort or limit, and streams rows back to the client. The leader also runs **leader-only queries**: statements that reference only catalog tables (the `PG_*` catalog, some `STV_`/`SVV_` views) or leader-only functions never reach the compute nodes. This is the source of a classic Redshift error when someone uses a leader-only function against a user table in the same query — the plan cannot be split, and Redshift refuses it. ## What compute nodes do Each compute node is subdivided into **slices**, and each slice is an independent worker with its own share of the node's memory and disk. A table's rows are spread across all the slices in the cluster according to its distribution style, and every slice runs the same compiled segment over the rows it owns. Compute nodes: - scan, filter, aggregate, join and sort their own data; - exchange rows directly with each other over the interconnect when a step requires redistribution or broadcast; - send only their partial results to the leader. Because slices work independently, the wall-clock time of a step is set by the *slowest* slice, not the average one — which is why uneven data placement hurts so much in this architecture. ## Single-node clusters On a single-node cluster, the one node plays both roles: it is leader and compute simultaneously. This is fine for development and tiny datasets, but it means no cross-node parallelism, and it is not how production clusters are shaped. ## RA3 does not change the split On RA3 node types the authoritative copy of the data lives in Redshift Managed Storage rather than only on node-local disk, and local SSD becomes a cache. The leader/compute division of labour is exactly the same: the leader still plans and merges, compute nodes still own slices and execute segments. ## Why the split shows up in everyday SQL Several common performance surprises trace directly back to this architecture: - **Huge result sets.** Every returned row passes through the single leader node. `SELECT * FROM fact_events;` against a billion rows makes a 16-node cluster behave like a one-machine funnel. Aggregate on the cluster, or use `UNLOAD`, which writes to S3 in parallel straight from the slices without routing rows through the leader. - **Final sorts.** An `ORDER BY` at the top of a plan is merged on the leader; with a `LIMIT` that is cheap, without one it is not. - **Connection pressure.** Sessions and their state live on the leader, and a cluster has a finite connection limit, so hundreds of idle pooled connections are a leader-side cost, not a compute-side one. - **First-run latency.** A dashboard whose SQL text changes on every run (inlined literals instead of parameters) keeps missing the compiled-code cache. ## What an interviewer is listening for That you can say *one leader, many compute nodes, slices underneath the compute nodes*, that the leader holds no user data, and that you can name at least one query shape whose cost is explained by the leader being a single point of serialization.

  • Why can returning a very large result set be slow even on a large Redshift cluster?
    Every row travels back through the single leader node, which merges the compute nodes' partial results and streams them to the client. That step is not parallel, so cluster size does not help. Aggregate or sample on the cluster, or use UNLOAD to S3, which writes in parallel directly from the slices and bypasses the leader entirely.
  • Why is the first execution of a new query shape sometimes much slower than the second?
    Redshift's leader node compiles the plan's segments into executable code before shipping them to the compute nodes. That compilation happens once per query shape and the result is cached, so repeats are fast. Queries that inline literals instead of using parameters change shape every run and keep missing the cache.
  • What happens to the leader/compute split on a single-node Redshift cluster?
    The single node performs both roles: it plans, compiles and merges like a leader, and stores and scans data like a compute node. It is useful for development, but there is no cross-node parallelism and no isolation between coordination work and scan work.

The leader is a conductor: it reads the score, hands each section its part, and cues the final chord — but it never plays an instrument or holds any of the sheet music the orchestra is reading from.

saying these in an interview costs you the question

  • Claiming the leader node stores a copy of the data
  • Saying compute nodes plan their own queries independently
  • Thinking clients connect directly to compute nodes
  • Assuming a bigger cluster speeds up returning millions of rows
  • Confusing the leader node with a WLM queue or a coordinator process per query

context