Amazon Redshift is a SQL data warehouse, and most people meet it as a JDBC endpoint that accepts PostgreSQL-flavoured queries. That surface hides the part that decides whether a query takes two seconds or two hours: Redshift is a massively parallel, columnar engine, and your table definitions tell it how to spread data across many workers. Get that layout right and joins over billions of rows stay local to each worker; get it wrong and every query shuffles the whole table across the network.

This article explains the machine from first principles, then works through a fact and dimension design, a load pipeline, workload management and the system views you read when something is slow. If you are still deciding between a warehouse and querying files in place, the Athena article and the data lake article cover the other side of that choice.

Advertisement

The machine: a leader, compute nodes and slices

A provisioned Redshift cluster has one leader node and several compute nodes. The leader node is the only thing clients talk to. It parses SQL, plans with table statistics, compiles the plan, sends it to the compute nodes and merges their results.

Each compute node is divided into slices. A slice owns a portion of the node's memory and disk, and every table's rows are assigned to slices. When a query runs, every slice scans and processes its own rows in parallel, so a query's speed is set by the slowest slice. That single sentence explains most Redshift tuning: you want each slice to hold a similar amount of data, and you want rows that are joined together to sit on the same slice. Query SELECT node, slice FROM stv_slices to see how many slices your cluster has instead of relying on a published table.

Current RA3 node types separate compute from storage. Data lives in Redshift Managed Storage, which is backed by S3, and each node caches hot blocks on local SSD. Redshift Serverless goes one step further: you set a base capacity in Redshift Processing Units (RPUs), the service scales above it under load, and you pay for RPU-hours used. AWS documents a default base of 128 RPUs, a configurable range of 4 to 512, and 16 GB of memory per RPU. The execution model underneath is the same, so everything below applies to both.

Redshift query path: one leader plans, every slice scans its share of the columnsSQL clientsBI, dbt, JDBC/ODBCLeader nodeparse, plan, compile, mergeWLM queuespriority, SQA, scalingqueryadmitCompute node 1slice 0 | slice 1Compute node 2slice 2 | slice 3Compute node 3slice 4 | slice 5segmentsredistribute or broadcast rows between slices only when join keys are not co-locatedRedshift Managed Storage (RA3 / Serverless)columnar 1 MB blocks in S3-backed storage; hot blocks cached on local SSD; zone maps hold min/max per blockAmazon S3 data lakeCOPY / UNLOAD, Spectrum external tablesOther warehousesdata sharing: live reads, no copiesServerless replaces the node count with RPUs, but the plan, slices, blocks and data movement work the same way.
Clients reach only the leader. The leader compiles a plan and fans it out to every slice; slices exchange rows only when the data layout forces them to.

Columnar blocks, compression and zone maps

Redshift stores each column separately in 1 MB blocks. A query that reads three columns of a forty-column table reads roughly three fortieths of the data, which is the first reason analytic queries are fast. Single-type blocks compress well: AZ64 suits numbers and timestamps, ZSTD general text, and RAW means none. Automatic encoding picks for you and leaves sort key columns RAW.

The second reason is zone maps. For every block, Redshift records the minimum and maximum value. When a query filters on a column, a slice skips any block whose range cannot match. Zone maps only help if values are clustered, which is what the sort key does. A table sorted by order_ts has blocks that each cover a narrow time window, so a query for one month touches only that month's blocks. For the general theory of column stores, see the columnar storage article.

Advertisement

Distribution styles: deciding which slice owns a row

Every table has a distribution style, and it is the most important decision in a Redshift design because it decides whether joins move data.

StyleWhere rows goUse it forWatch out for
KEYHash of one column decides the sliceLarge tables joined to each other on that columnSkew when a few key values dominate, for example a null or a test customer
ALLA full copy on every nodeSmall, slowly changing dimensionsLoad cost and storage multiply by the node count
EVENRound-robin across slicesTables with no dominant join, or staging tablesJoins on such tables always need redistribution or broadcast
AUTOStarts as ALL for small tables, moves to EVEN or KEY as it growsThe default when you have no evidence yetIt optimises from observed workload, which takes time to learn

When two tables are distributed on the join column, matching rows already sit on the same slice and the join runs locally. When they are not, the planner either broadcasts the smaller side to every node or redistributes one or both sides by the join key. You see the decision in EXPLAIN output, and it is the first thing to check on any slow join.

EXPLAIN
SELECT c.region, date_trunc('day', o.order_ts) AS day, sum(o.amount)
FROM   sales.orders o
JOIN   sales.customers c ON c.customer_id = o.customer_id
WHERE  o.order_ts >= '2026-09-01' AND o.order_ts < '2026-10-01'
GROUP  BY 1, 2;

-- Look for the join step's distribution label:
--   DS_DIST_NONE / DS_DIST_ALL_NONE   no rows move: co-located or the inner is ALL
--   DS_BCAST_INNER                    the inner table is broadcast to every node
--   DS_DIST_INNER / DS_DIST_OUTER     one side is redistributed on the join key
--   DS_DIST_BOTH                      both sides move: almost always a design bug

Sort keys and the unsorted region

A compound sort key orders rows by the first column, then the second within ties, and so on. It helps filters and merge joins that use a prefix of the key, and it does little for filters on later columns alone. Interleaved keys cost more to maintain, so prefer compound or AUTO. The usual first choice is the column most queries filter by range, which for event data is a timestamp.

New rows land in an unsorted region at the end of the table. Automatic vacuum sorts and purges in the background, but heavy out-of-order loads, updates and deletes can outpace it, and zone maps lose their power. Run VACUUM SORT ONLY or VACUUM DELETE ONLY on a specific table during a quiet window, and load data in sort-key order where you can.

Worked example: an orders fact and a customer dimension

Suppose an online shop stores 4 billion orders and 20 million customers. Analysts ask for revenue by region and day over recent months, and another team joins orders to a 3 billion row clickstream table on customer id.

-- Dimension: small, joined from everywhere -> copy it to every node.
CREATE TABLE sales.customers (
    customer_id   BIGINT       NOT NULL,
    region        VARCHAR(16)  NOT NULL,
    segment       VARCHAR(16)  NOT NULL,
    created_at    TIMESTAMP    NOT NULL,
    PRIMARY KEY (customer_id)            -- informational only: NOT enforced
)
DISTSTYLE ALL
SORTKEY (customer_id);

-- Fact: large, joined to a large table on customer_id, filtered by time.
CREATE TABLE sales.orders (
    order_id      BIGINT         NOT NULL,
    customer_id   BIGINT         NOT NULL,
    order_ts      TIMESTAMP      NOT NULL,
    status        VARCHAR(16)    NOT NULL ENCODE zstd,
    amount        DECIMAL(12,2)  NOT NULL ENCODE az64,
    PRIMARY KEY (order_id)
)
DISTSTYLE KEY
DISTKEY (customer_id)
COMPOUND SORTKEY (order_ts);

The customer table is distributed ALL, so every node has a full copy and the join from orders never moves rows; at 20 million narrow rows the replication is affordable. Orders is distributed by customer id because the other large table it joins with also uses customer id; distributing the clickstream on the same key makes that big-to-big join local. Orders is sorted by timestamp because the dominant filter is a date range, so a one-month query scans roughly one forty-eighth of a four-year table.

Check the choice against data before shipping it. If one customer id accounts for a large share of orders, for example a marketplace account that places orders on behalf of others, the slice owning that key becomes a hotspot and every query waits for it. svv_table_info.skew_rows reports the ratio between the fullest and emptiest slice; values well above 1 deserve a look. The fix is a higher-cardinality key, accepting data movement for the rarer join.

Note the primary keys. Redshift accepts them and the planner may use them, but it does not enforce them. Loading the same file twice silently doubles revenue, so enforce uniqueness in the load itself.

Loading: COPY, staging and MERGE

Single-row INSERT statements are the slowest way to put data into Redshift: each one is a small transaction, and commits are serialised across the cluster. Load in bulk with COPY from S3, which reads files in parallel across slices. Split input into several files, ideally a multiple of the slice count and of roughly equal size; AWS guidance is files between about 1 MB and 1 GB after compression. Columnar formats such as Parquet need no delimiter or quoting options. The S3 layout and permissions side is covered in the S3 article, and building the files with Glue jobs in the Glue ETL article.

-- 1. Land the batch in a staging table with the same layout as the target.
CREATE TEMP TABLE stage_orders (LIKE sales.orders);

-- 2. COPY reads every file under the prefix in parallel, one file per slice at a time.
COPY stage_orders
FROM 's3://acme-lake/orders/dt=2026-10-01/'
IAM_ROLE 'arn:aws:iam::123456789012:role/redshift-copy'
FORMAT AS PARQUET;

-- 3. Upsert. MERGE fails if one target row matches several source rows, so keep
--    exactly one row per order (the newest), which also drops exact duplicates.
CREATE TEMP TABLE stage_dedup AS
SELECT order_id, customer_id, order_ts, status, amount
FROM (SELECT *, ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY order_ts DESC) AS rn
      FROM stage_orders)
WHERE rn = 1;

MERGE INTO sales.orders
USING stage_dedup s ON sales.orders.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET status = s.status, amount = s.amount, order_ts = s.order_ts
WHEN NOT MATCHED THEN INSERT (order_id, customer_id, order_ts, status, amount)
                      VALUES (s.order_id, s.customer_id, s.order_ts, s.status, s.amount);

-- 4. Keep the planner honest after a large change.
ANALYZE sales.orders;

This pattern makes the load idempotent: rerunning it for the same day replaces rows instead of duplicating them. Wrap the steps in one transaction.

Workload management and concurrency

Workload management (WLM) decides which queries run now and how much memory each gets. Automatic WLM, the default, estimates memory per query and adjusts concurrency, and you steer it with queues and priorities: route dashboards and ETL to separate queues by user group or query group, and give interactive work higher priority. Short query acceleration runs queries predicted to be short in a dedicated space so they do not wait behind long scans. Query monitoring rules can log, deprioritise or abort runaway queries.

When queued work exceeds the cluster, concurrency scaling can add transient capacity for eligible queries, billed per second beyond a daily free allowance; check current pricing before relying on it. Serverless behaves similarly through automatic RPU scaling, so set a maximum RPU-hours usage limit to bound cost. Neither makes a single badly laid-out query faster.

Beyond the cluster: Spectrum, materialized views and data sharing

Redshift Spectrum queries external tables over files in S3, defined in the Glue Data Catalog, and joins them with local tables. It suits cold or rarely queried history: keep recent months local and leave older partitions in the lake.

Materialized views precompute expensive aggregates and can refresh automatically. Data sharing, available on RA3 and Serverless, lets a consumer warehouse query a producer's live data without copying it. Zero-ETL integrations replicate operational databases, such as Aurora, into Redshift without a pipeline you operate.

Failure modes and how to diagnose them

SymptomLikely causeWhat to check or do
One query is slow, cluster otherwise idleData movement: DS_BCAST_INNER of a big table or DS_DIST_BOTHRead EXPLAIN; align distribution keys of the large tables that join
All queries slower over weeksUnsorted blocks and deleted rows, stale statisticssvv_table_info unsorted and stats_off; VACUUM and ANALYZE large tables
Some queries far slower on one daySkew: one slice holds most rowsskew_rows; choose a higher-cardinality distribution key
Queries spill to diskToo little memory per query, or huge intermediate resultsis_diskbased in svl_query_summary; reduce columns, pre-aggregate, adjust WLM
Totals doubled after a rerunUnenforced primary keys plus a non-idempotent loadStage, deduplicate and MERGE; add a uniqueness check to the pipeline

Three queries cover most investigations: table health, the planner's alerts, and the slowest recent queries.

-- Table health: skew, unsorted fraction and stale statistics.
SELECT "table", diststyle, sortkey1, skew_rows, unsorted, stats_off, tbl_rows, size
FROM   svv_table_info
ORDER  BY size DESC
LIMIT  20;

-- The planner's own warnings.
SELECT query, trim(event) AS event, trim(solution) AS solution
FROM   stl_alert_event_log
WHERE  event_time > dateadd(day, -1, getdate());

-- Slowest recent queries (SYS_ views work on provisioned and Serverless).
SELECT query_id, status, queue_time, execution_time, elapsed_time, left(query_text, 80)
FROM   sys_query_history
ORDER  BY elapsed_time DESC
LIMIT  20;

Trade-offs: when Redshift is the right tool

Redshift fits repeated analytic queries over large, well-modelled data with many concurrent users and predictable latency. Athena fits occasional queries over files you do not want to load. Provisioned RA3 suits steady utilisation; Serverless suits spiky use. Redshift is a poor fit for point lookups, high-rate single-row writes and transactional workloads, which belong in an OLTP database such as DynamoDB or Aurora.

What to do next

  1. List your five most frequent and five slowest queries, and note their join columns and filter columns.
  2. Run the svv_table_info query and record skew_rows, unsorted and stats_off for your largest tables.
  3. Run EXPLAIN on each slow query and find any DS_BCAST_INNER on a large table or DS_DIST_BOTH; realign distribution keys for those tables.
  4. Pick a compound sort key that matches the dominant range filter, and load data in that order.
  5. Replace row-by-row inserts with COPY into staging plus deduplicate and MERGE, and make each load idempotent.
  6. Separate ETL and interactive work into WLM queues with priorities, add query monitoring rules, and set usage limits on concurrency scaling or Serverless RPU-hours.
Key takeaway: Redshift is a parallel, columnar engine in which every slice works on its own rows, so performance comes from layout. Distribute large tables on the column they join by, copy small dimensions to every node, sort by the dominant range filter, and load in bulk through staging and MERGE because keys are not enforced. When a query is slow, read EXPLAIN for data movement and svv_table_info for skew, unsorted blocks and stale statistics before buying more capacity.