AWS courseLesson 7 of 12
AWS course · Lesson 7 of 12
Amazon Redshift for Data Engineers
Redshift for Data Engineers: architecture, RA3, distribution and sort keys, COPY and UNLOAD, vacuum, WLM, concurrency scaling, Spectrum and Serverless.
On this page
Amazon Redshift is AWS’s data warehouse: a columnar, massively parallel SQL database built for analytical queries over large fact tables, with many concurrent BI users. It looks like PostgreSQL from a SQL client, but it stores and processes data very differently, and most Redshift interview questions test that difference. This lesson covers how Redshift spreads data across nodes, how to load it, how to keep it fast, and when to use Serverless.
Sample data
Redshift needs an AWS account, so its SQL is shown but not executed. Two concepts get runnable stand-ins: a Python model of how distribution keys spread rows across slices, and PostgreSQL 16 analogues for zone maps (a BRIN index) and materialized views. They are labelled as analogues: PostgreSQL is row-oriented and single-node, so it shows the idea, not Redshift’s behaviour.
-- PostgreSQL analogue table: 500,000 time-ordered events
CREATE TABLE events (event_id bigint, event_ts timestamp, user_id int, amount numeric(10,2));
INSERT INTO events
SELECT g, timestamp '2026-01-01' + (g * interval '30 seconds'), g % 5000, (g % 100) / 4.0
FROM generate_series(1, 500000) AS g;
ANALYZE events;
Redshift architecture
What it is. A provisioned Redshift cluster has one leader node and one or more compute nodes. Each compute node is split into slices, each with a share of the node’s CPU, memory and data.
How it works.
- A client connects to the leader node (PostgreSQL wire protocol, JDBC/ODBC, or the Data API).
- The leader parses the query, builds a plan, compiles it into code, and sends segments to the compute nodes.
- Each slice scans and processes its own rows in parallel. Rows that must meet for a join or aggregation are moved between nodes over the network (redistribution or broadcast).
- The leader combines the partial results and returns them.
Storage is columnar: each column is stored in 1 MB blocks, compressed with a column encoding, and each block has min/max metadata called a zone map, which lets scans skip blocks that cannot match a filter. Redshift has no indexes in the usual sense; zone maps plus sort keys do that job.
| Concept | Redshift | Typical OLTP PostgreSQL |
|---|---|---|
| Storage | Columnar, compressed 1 MB blocks | Row-oriented pages |
| Parallelism | Many slices on many nodes | Mostly one server |
| Skipping data | Zone maps, sort keys | B-tree indexes |
| Constraints | Primary and foreign keys are informational, not enforced | Enforced |
| Best at | Scans, joins and aggregations over billions of rows | Many small reads and writes |
Pitfalls.
- Treating Redshift like an OLTP database: single-row
INSERTs in a loop, frequentUPDATEs of individual rows, or relying on unenforced unique keys for de-duplication. - Functions that only run on the leader node (some system catalogue functions) fail when mixed with table data.
In interviews. Draw the leader node, compute nodes and slices, and explain that query speed depends on how much data each slice scans and how much moves across the network. That sets up distribution and sort keys.
RA3 and managed storage
What it is. RA3 nodes separate compute from storage. Data lives in Redshift managed storage (RMS), backed by S3, while each node keeps hot data in a large local SSD cache. You size the cluster for compute and pay for storage separately.
How it works. RA3 node types are ra3.large, ra3.xlplus, ra3.4xlarge and ra3.16xlarge. For example, ra3.xlplus has 4 vCPUs, 32 GiB of memory and a managed storage quota of 32 TB per node; ra3.4xlarge has 12 vCPUs and 96 GiB; ra3.16xlarge has 48 vCPUs and 384 GiB, both with 128 TB managed storage quota per node. Older DC2 nodes keep all data on local SSDs, so storage and compute scale together.
Separating storage enables several features that matter to data engineers:
- Data sharing: other clusters and Serverless workgroups (in the same or other accounts) query live data without copying it.
- Elastic resize and pause and resume without moving all data.
- Concurrency scaling for writes, and Serverless itself, which uses the same managed storage.
Pitfalls. Sizing an RA3 cluster by data volume, the DC2 habit. Size by compute needs (query concurrency and complexity) and let managed storage grow.
In interviews. Explain RA3 as “compute nodes with a local cache over S3-backed managed storage”, and link it to data sharing and to Serverless. Compare it with Snowflake’s separation of storage and compute if asked.
Distribution styles
What it is. The distribution style decides which slice stores each row. Good distribution keeps related rows on the same slice so joins happen locally, and spreads rows evenly so no slice does more work than the others.
How it works.
| Style | Rows go to | Use for |
|---|---|---|
AUTO (default) |
Redshift chooses and changes the style as the table grows | Most tables, especially when unsure |
EVEN |
Round-robin across slices | Tables not joined on a consistent key |
KEY |
Slice chosen by hashing one column | Large fact table and its largest dimension joined on that column |
ALL |
A full copy on every node | Small, slowly changing dimension tables |
CREATE TABLE sales_fact (
sale_id bigint NOT NULL,
customer_id integer NOT NULL,
product_id integer NOT NULL,
sale_ts timestamp NOT NULL,
amount decimal(12,2)
)
DISTSTYLE KEY
DISTKEY (customer_id)
SORTKEY (sale_ts);
CREATE TABLE dim_customer (customer_id integer, segment varchar(20), country char(2))
DISTSTYLE KEY DISTKEY (customer_id);
CREATE TABLE dim_currency (code char(3), name varchar(40)) DISTSTYLE ALL;
The distribution key must have many distinct, evenly used values. This model hashes 100,000 orders onto 8 slices by two candidate keys: customer_id (well spread) and country (40% of rows share one value):
import hashlib
from collections import Counter
SLICES = 8
def slice_for(value):
"""Stand-in for Redshift's hash distribution: the same key always lands on the same slice."""
return int(hashlib.md5(str(value).encode()).hexdigest(), 16) % SLICES
# 100,000 orders: customer_id is well spread, but 40% of rows have the same country.
orders = [{"order_id": i, "customer_id": i % 9973, "country": "GB" if i % 5 < 2 else f"C{i % 37}"}
for i in range(100_000)]
def skew(column):
counts = Counter(slice_for(o[column]) for o in orders)
rows = [counts.get(s, 0) for s in range(SLICES)]
return max(rows) / (sum(rows) / SLICES), rows
for col in ("customer_id", "country"):
ratio, rows = skew(col)
print(f"DISTKEY({col:<11}) max/avg rows per slice = {ratio:.2f} {rows}")
even = [len(range(s, 100_000, SLICES)) for s in range(SLICES)]
print(f"DISTSTYLE EVEN max/avg rows per slice = {max(even) / (100_000 / SLICES):.2f} {even}")
DISTKEY(customer_id) max/avg rows per slice = 1.04 [12650, 12006, 12240, 12507, 12558, 12568, 12500, 12971]
DISTKEY(country ) max/avg rows per slice = 3.85 [6488, 6486, 9728, 48110, 6488, 6486, 8106, 8108]
DISTSTYLE EVEN max/avg rows per slice = 1.00 [12500, 12500, 12500, 12500, 12500, 12500, 12500, 12500]
With DISTKEY(country) one slice holds almost four times the average, so every query on that table waits for that slice. In Redshift, check skew in SVV_TABLE_INFO (skew_rows) and look for DS_DIST_* and DS_BCAST_* steps in EXPLAIN output, which show rows being redistributed or broadcast.
Pitfalls.
- Low-cardinality or skewed keys (country, status, a date) cause hot slices.
- Distributing the fact table on a key that is not its main join key gives no co-location benefit.
ALLon a large or frequently updated table multiplies storage and slows loads.
In interviews. Explain co-located joins (fact and largest dimension on the same DISTKEY), ALL for small dimensions, EVEN when there is no dominant join, and AUTO as the sensible default. Name skew as the main risk.
Sort keys
What it is. A sort key orders rows on disk within each slice. Combined with zone maps, it lets range filters skip most blocks: if data is sorted by sale_ts, a one-day query reads only the blocks for that day.
How it works.
- Compound sort key (the default kind): sorts by the first column, then the second, and so on. Best when filters use the leading column, typically a timestamp.
- Interleaved sort key: gives equal weight to each column; useful for very specific multi-column filter patterns, but costly to maintain and rarely the right choice now.
SORTKEY AUTO: Redshift picks and adjusts sort keys from query history.
Newly loaded rows go into an unsorted region until vacuumed or sorted automatically, which reduces skipping.
The closest PostgreSQL analogue to zone maps is a BRIN index, which stores min/max values per range of pages and works best when the table is physically ordered by the column. Sequential scans are disabled here only to show the BRIN plan on this small table:
CREATE INDEX events_ts_brin ON events USING brin (event_ts);
SET enable_seqscan = off;
EXPLAIN (COSTS OFF) SELECT count(*) FROM events WHERE event_ts >= '2026-03-01' AND event_ts < '2026-03-02';
RESET enable_seqscan;
SELECT count(*) FROM events WHERE event_ts >= '2026-03-01' AND event_ts < '2026-03-02';
Aggregate
-> Bitmap Heap Scan on events
Recheck Cond: ((event_ts >= '2026-03-01 00:00:00'::timestamp without time zone) AND (event_ts < '2026-03-02 00:00:00'::timestamp without time zone))
-> Bitmap Index Scan on events_ts_brin
Index Cond: ((event_ts >= '2026-03-01 00:00:00'::timestamp without time zone) AND (event_ts < '2026-03-02 00:00:00'::timestamp without time zone))
count
-------
2880
Pitfalls.
- Sorting by a column nobody filters on.
- Wrapping the sort column in a function in the
WHEREclause (WHERE date_trunc('day', sale_ts) = ...), which can prevent block skipping. Use a range on the raw column. - Compressing the leading sort key column too aggressively; Redshift recommends leaving the first sort key column raw-encoded (or letting
AUTOdecide), because heavy compression of it can make zone maps less selective.
In interviews. “Sort key for time-range filters, distribution key for joins” is the core sentence. Add zone maps as the mechanism and the unsorted region as the reason vacuum matters.
COPY and UNLOAD
What it is. COPY is the fast, parallel way to load data from S3 (and other sources) into Redshift. UNLOAD writes query results back to S3 in parallel.
How it works. COPY splits the work across slices, so it is fastest when the input is many files of similar size (a multiple of the number of slices) or columnar Parquet. Authenticate with an IAM role attached to the cluster or namespace, never with access keys.
COPY sales_fact
FROM 's3://example-lake/curated/sales/dt=2026-10-03/'
IAM_ROLE 'arn:aws:iam::111122223333:role/redshift-copy'
FORMAT AS PARQUET;
COPY staging_orders
FROM 's3://example-lake/raw/orders/2026-10-03/'
IAM_ROLE 'arn:aws:iam::111122223333:role/redshift-copy'
FORMAT AS CSV IGNOREHEADER 1 GZIP
TIMEFORMAT 'auto'
MAXERROR 0;
UNLOAD ('SELECT customer_id, sum(amount) AS revenue FROM sales_fact GROUP BY customer_id')
TO 's3://example-exports/customer_revenue/'
IAM_ROLE 'arn:aws:iam::111122223333:role/redshift-unload'
FORMAT AS PARQUET
PARTITION BY (customer_id)
MAXFILESIZE 256 MB;
A common load pattern is staging and merge: COPY into a staging table, then MERGE (or DELETE plus INSERT in one transaction) into the target, so reruns do not duplicate rows. Load errors are recorded in SYS_LOAD_ERROR_DETAIL (or the older STL_LOAD_ERRORS).
Pitfalls.
- One huge gzip file: it cannot be split, so one slice does all the loading.
INSERT ... VALUESrow by row from an application: very slow; batch through S3 andCOPY.- Re-running a
COPYof the same files duplicates data. Use a staging table and merge, or a manifest that lists exactly which files to load.
In interviews. Explain why COPY beats inserts (parallel per slice), how to make loads idempotent (staging plus merge, manifests), and that UNLOAD with Parquet and PARTITION BY feeds the lake.
Vacuum and analyze
What it is. Redshift does not overwrite rows in place. An UPDATE marks the old row deleted and writes a new one; a DELETE only marks rows. VACUUM reclaims the space from deleted rows and re-sorts rows into sort key order. ANALYZE refreshes the statistics the planner uses.
How it works. Redshift runs automatic vacuum delete, automatic table sort and automatic analyze in the background when the cluster is lightly loaded, so many teams rarely run these by hand. After large deletes or loads that land out of order you can run them explicitly:
VACUUM DELETE ONLY sales_fact; -- reclaim space only
VACUUM SORT ONLY sales_fact; -- re-sort without reclaiming
VACUUM FULL sales_fact TO 99 PERCENT; -- both, until 99% sorted
ANALYZE sales_fact PREDICATE COLUMNS; -- statistics for columns used in filters and joins
SELECT "table", unsorted, stats_off, tbl_rows
FROM svv_table_info
ORDER BY unsorted DESC NULLS LAST;
PostgreSQL has commands with the same names, but they differ: PostgreSQL’s VACUUM marks dead row versions reusable and does not re-sort the table, and its ANALYZE samples rows for a single-node planner. The ideas they share are that deletes leave garbage behind and that planners need fresh statistics.
Pitfalls.
- Frequent small
UPDATEs andDELETEs, which create many deleted rows and unsorted blocks. Prefer batch merges. - Running
VACUUM FULLon huge tables during business hours; it is resource-heavy. Let automatic vacuum work, or schedule it. - Deep copy (
CREATE TABLE ASand swap) is sometimes faster than vacuuming a badly unsorted table.
In interviews. Explain why deletes and updates leave work behind in a columnar MPP store, what VACUUM and ANALYZE each fix, and that Redshift now automates most of it.
Workload management (WLM)
What it is. WLM decides how many queries run at once and how memory is shared, so heavy ETL does not starve dashboards.
How it works.
- Automatic WLM (the default and recommended mode) lets Redshift choose concurrency and memory per query. You define queues mapped to user groups or query groups and give them priorities (highest to lowest).
- Manual WLM sets a fixed number of slots and memory percentage per queue.
- Query monitoring rules (QMR) act on runaway queries, for example log, change priority or abort a query that scans too many rows or runs too long.
- Short query acceleration (SQA) fast-tracks short queries so they do not wait behind long ones.
Redshift Serverless replaces queues with workgroup settings and query limits, but the same idea applies: separate and cap workloads.
Pitfalls.
- One queue for everything, so a nightly backfill blocks the morning dashboards.
- Manual WLM with too many slots, which gives each query too little memory and makes them spill to disk.
In interviews. Describe separate queues (ETL, BI, ad hoc) with priorities, query monitoring rules to stop runaway queries, and SQA for short queries.
Concurrency scaling
What it is. When queries start queueing, concurrency scaling adds temporary clusters that run the queued queries, then removes them. Users see consistent performance during bursts without a permanently bigger cluster.
How it works. You enable it per WLM queue. Read queries are eligible, and on RA3 many write operations are too (for example COPY, INSERT, DELETE, UPDATE and CTAS). Each active cluster earns up to one hour of free concurrency scaling credits per day, accumulating up to 30 hours; usage beyond the credits is billed per second at on-demand rates, with a one-minute minimum each time a scaling cluster activates. A usage limit caps the time to control cost.
Pitfalls.
- Turning it on for every queue without a usage limit.
- Expecting it to speed up a single slow query. It adds parallel capacity for more queries, not more power for one.
In interviews. Contrast concurrency scaling (more clusters for more simultaneous queries) with resizing (a bigger cluster for faster individual queries), and mention the free daily credits and usage limits.
Redshift Spectrum
What it is. Spectrum lets a Redshift cluster query data in S3 through external tables defined in the Glue Data Catalog, without loading it. A fleet of Spectrum nodes does the scanning and filtering, and Redshift joins the result with local tables.
How it works.
CREATE EXTERNAL SCHEMA lake
FROM DATA CATALOG
DATABASE 'curated'
IAM_ROLE 'arn:aws:iam::111122223333:role/redshift-spectrum';
SELECT c.segment, sum(s.amount) AS revenue
FROM lake.sales AS s -- Parquet in S3, partitioned by dt
JOIN dim_customer AS c ON c.customer_id = s.customer_id
WHERE s.dt BETWEEN '2026-09-01' AND '2026-09-30'
GROUP BY c.segment;
Spectrum is billed by data scanned in S3 (on provisioned clusters), so the same layout rules as Athena apply: Parquet, partitions, filter on partition columns. On Serverless, queries on S3 data are part of RPU usage.
Pitfalls.
- Joining a huge unpartitioned external table to local tables: the scan dominates. Partition and filter.
- Expecting partition projection to work: it is Athena-only. Spectrum needs partitions in the catalogue.
- The IAM role needs both S3 and Glue Data Catalog permissions (and Lake Formation grants if the data is governed).
In interviews. Use Spectrum for “keep recent hot data in Redshift and years of history in S3, and query both together”. Compare it with Athena: same catalogue and layout rules, but joins with warehouse tables and the warehouse’s concurrency controls.
Materialized views
What it is. A materialized view stores the result of a query and refreshes it later, so dashboards read a small precomputed table instead of re-aggregating a large fact table each time.
How it works. Redshift can refresh incrementally for many query shapes (it applies only the changes since the last refresh) and fully recompute others. AUTO REFRESH YES refreshes in the background when base tables change, and automatic query rewrite can make queries against the base tables use a matching materialized view.
CREATE MATERIALIZED VIEW mv_daily_revenue
AUTO REFRESH YES
AS
SELECT trunc(sale_ts) AS day, sum(amount) AS revenue, count(*) AS sales
FROM sales_fact
GROUP BY trunc(sale_ts);
REFRESH MATERIALIZED VIEW mv_daily_revenue;
PostgreSQL has materialized views too, but it only supports full refresh (no incremental refresh, no auto refresh). The analogue shows the key behaviour both share: the view is stale until refreshed.
CREATE MATERIALIZED VIEW daily_revenue AS
SELECT date_trunc('day', event_ts)::date AS day, sum(amount) AS revenue, count(*) AS events
FROM events GROUP BY 1;
INSERT INTO events VALUES (500001, '2026-01-01 12:00', 1, 1000.00);
SELECT * FROM daily_revenue WHERE day = '2026-01-01';
REFRESH MATERIALIZED VIEW daily_revenue;
SELECT * FROM daily_revenue WHERE day = '2026-01-01';
day | revenue | events
------------+----------+--------
2026-01-01 | 35440.00 | 2879
day | revenue | events
------------+----------+--------
2026-01-01 | 36440.00 | 2880
Pitfalls.
- Assuming freshness: auto refresh is asynchronous and lower priority than user queries.
- Query shapes that cannot refresh incrementally (some outer joins, window functions, certain functions) fall back to full recompute, which may be expensive on every refresh.
In interviews. Explain materialized views as precomputed aggregates for dashboards, the difference between incremental and full refresh, and the staleness trade-off.
Redshift Serverless
What it is. Redshift Serverless runs the warehouse without clusters. A namespace holds the database objects and data; a workgroup holds the compute settings and endpoint. Capacity is measured in Redshift Processing Units (RPUs) and scales automatically with load.
How it works.
- Base capacity is the RPU level that queries start with. The default is 128 RPUs; it can be set from as low as 4 RPUs up to 512, or 1,024 in some Regions.
- Maximum capacity (an RPU-hour limit) and usage limits cap spending, with actions such as alerting or turning off user queries.
- Compute is billed in RPU-hours, per second, with a 60-second minimum whenever queries run; there is no charge when idle. Managed storage is billed separately.
- AI-driven scaling options let you set a price-performance target instead of a fixed base.
aws redshift-serverless create-namespace --namespace-name analytics \
--iam-roles arn:aws:iam::111122223333:role/redshift-copy
aws redshift-serverless create-workgroup --workgroup-name analytics-wg \
--namespace-name analytics --base-capacity 32 --max-capacity 256
| Choose | When |
|---|---|
| Serverless | Intermittent or unpredictable workloads, new projects, teams that do not want to manage clusters |
| Provisioned RA3 | Steady, heavy usage around the clock, where reserved nodes cost less, or features and tuning not available in Serverless |
Pitfalls.
- Leaving the base capacity at 128 RPUs for a small workload, or not setting a maximum, so a bad query scales up the bill.
- Keeping a BI tool polling constantly, which keeps compute running; each run is billed at least 60 seconds.
In interviews. Define RPUs, base and maximum capacity, per-second billing with a 60-second minimum, and namespaces versus workgroups. Then compare with provisioned clusters by workload shape.
Practice questions
How would you choose distribution and sort keys for a sales fact table joined mostly to customers and filtered by date?
Distribute the fact table and the customer dimension on customer_id (if it has high cardinality and even spread), so the main join is co-located; use ALL for small dimensions such as currency; sort the fact table by sale_ts (compound) so date filters skip blocks through zone maps. Check skew in SVV_TABLE_INFO and redistribution steps in EXPLAIN, or start with AUTO and let Redshift adjust.
A COPY of one 50 GB gzip CSV file is slow. Why, and how do you fix it?
Gzip files are not splittable, so a single slice loads the whole file while the others wait. Split the input into many similar-sized files (a multiple of the slice count), or convert it to Parquet, then COPY from the prefix or a manifest.
Dashboards slow down every morning while ETL runs. What Redshift features help?
Separate workloads with WLM queues and priorities (BI higher than ETL), add query monitoring rules to stop runaway queries, enable short query acceleration, enable concurrency scaling with a usage limit for the BI queue, and precompute dashboard aggregates with materialized views. On Serverless, put ETL and BI in separate workgroups using data sharing.
When would you use Spectrum instead of loading data into Redshift?
For large, rarely queried history kept in S3, data shared with other engines through the Glue Data Catalog, or data that changes in the lake and should not be copied. Keep hot, frequently joined data in Redshift tables for performance, and query cold data through Spectrum with partitions and Parquet.
What does VACUUM do in Redshift, and why is it needed?
Updates and deletes only mark rows as deleted, and new rows land in an unsorted region. VACUUM reclaims the space of deleted rows and re-sorts data into sort key order, which keeps zone maps effective. Redshift runs automatic vacuum delete and sort in the background, but large deletes or out-of-order loads can still need an explicit vacuum or a deep copy. ANALYZE refreshes statistics for the planner.
How is Redshift Serverless billed, and how do you keep it under control?
Compute is billed in RPU-hours, per second, with a 60-second minimum, only while queries run, plus managed storage. Control it with an appropriate base capacity, a maximum capacity, usage limits with alerts or query shut-off, query limits, and by avoiding constant polling from BI tools.
Key takeaways
- Redshift is a columnar MPP warehouse: a leader node plans, compute node slices scan in parallel, and zone maps skip blocks.
- RA3 nodes separate compute from S3-backed managed storage, which enables data sharing and Serverless.
- Distribution keys co-locate joins and must avoid skew; sort keys make range filters skip data.
- Load with parallel
COPYfrom many files or Parquet, make loads idempotent with staging and merge, and export withUNLOAD. - WLM queues, priorities, query monitoring rules and concurrency scaling keep mixed workloads predictable.
- Serverless bills RPU-hours per second with a 60-second minimum; base capacity defaults to 128 RPUs and can start at 4.
Progress is saved in this browser only. No account needed.