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.
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.
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.
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.
| Style | Where rows go | Use it for | Watch out for |
|---|---|---|---|
| KEY | Hash of one column decides the slice | Large tables joined to each other on that column | Skew when a few key values dominate, for example a null or a test customer |
| ALL | A full copy on every node | Small, slowly changing dimensions | Load cost and storage multiply by the node count |
| EVEN | Round-robin across slices | Tables with no dominant join, or staging tables | Joins on such tables always need redistribution or broadcast |
| AUTO | Starts as ALL for small tables, moves to EVEN or KEY as it grows | The default when you have no evidence yet | It 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
| Symptom | Likely cause | What to check or do |
|---|---|---|
| One query is slow, cluster otherwise idle | Data movement: DS_BCAST_INNER of a big table or DS_DIST_BOTH | Read EXPLAIN; align distribution keys of the large tables that join |
| All queries slower over weeks | Unsorted blocks and deleted rows, stale statistics | svv_table_info unsorted and stats_off; VACUUM and ANALYZE large tables |
| Some queries far slower on one day | Skew: one slice holds most rows | skew_rows; choose a higher-cardinality distribution key |
| Queries spill to disk | Too little memory per query, or huge intermediate results | is_diskbased in svl_query_summary; reduce columns, pre-aggregate, adjust WLM |
| Totals doubled after a rerun | Unenforced primary keys plus a non-idempotent load | Stage, 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
- List your five most frequent and five slowest queries, and note their join columns and filter columns.
- Run the svv_table_info query and record skew_rows, unsorted and stats_off for your largest tables.
- 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.
- Pick a compound sort key that matches the dominant range filter, and load data in that order.
- Replace row-by-row inserts with COPY into staging plus deduplicate and MERGE, and make each load idempotent.
- 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.