How do you choose distribution keys for a Redshift schema whose fact table joins three dimensions on different columns?
answer
- Only one of these per table, so it is a budget
- Rank the joins by bytes moved, not by count
- Small tables can sidestep the problem entirely
- A skewed key taxes every scan, not just the join
- The choice is revisable in place, so measure and revisit
basics
~20 sA Redshift table has one DISTKEY, so a fact can co-locate only one join. Spend it on the join that moves the most bytes, replicate small dimensions with DISTSTYLE ALL so their joins are local anyway, and accept redistribution for the rest.
solid answer
~50 sStart from evidence, not intuition: find the queries that dominate cluster time and see which joins actually move data, reading the distribution labels in their plans. Then apply the constraint. One DISTKEY per table means the fact table can co-locate exactly one join. Rank the candidate joins by bytes moved — a join to a 5 GB dimension matters far more than one to a 5 MB dimension — and give the DISTKEY to the winner, provided that column is high-cardinality and evenly distributed. Every small dimension should be `DISTSTYLE ALL`, which makes its join local regardless of the fact's key and usually removes the second and third problems entirely. If no candidate column is both a hot join key and uniform, `EVEN` is the honest answer and you pay redistribution knowingly. Where you genuinely cannot characterise the workload, leave tables on `AUTO` and let Automatic Table Optimization propose keys from observed queries — then review its recommendations rather than trusting them blindly.
code
sql · 25 lines-- large fact: spend the single DISTKEY on the biggest join
CREATE TABLE fact_order_items (
order_id BIGINT,
product_id INT,
store_id INT,
qty INT,
order_ts TIMESTAMP
)
DISTSTYLE KEY DISTKEY (order_id) -- co-locates the fact-to-fact join
COMPOUND SORTKEY (order_ts); -- prunes the time filter every query carries
-- large sibling fact shares the key -> DS_DIST_NONE
CREATE TABLE fact_orders (
order_id BIGINT,
customer_id BIGINT,
order_ts TIMESTAMP
)
DISTSTYLE KEY DISTKEY (order_id)
COMPOUND SORTKEY (order_ts);
-- small dimensions: replicate, so their joins are local anyway
CREATE TABLE dim_product (product_id INT, name VARCHAR(200))
DISTSTYLE ALL SORTKEY (product_id);
CREATE TABLE dim_store (store_id INT, region VARCHAR(40))
DISTSTYLE ALL SORTKEY (store_id);go deeper
Know that a Redshift table has one distribution key and that small dimension tables are often replicated to every node instead.
Explain why only one join can be co-located per fact table, and how DISTSTYLE ALL removes the problem for small dimensions.
Choose from evidence: rank joins by bytes moved, veto skewed or low-cardinality keys, and verify the result in plans rather than in theory.
Own the schema-wide policy — where AUTO is allowed and where keys are pinned, the storage budget for replicated dimensions, the review step for new large tables, and the cadence for revisiting the design as the workload shifts.
## The constraint that shapes everything A Redshift table has exactly one distribution key. That single fact drives the whole design conversation for a star schema: a fact table joined to `dim_customer`, `dim_product` and `dim_store` on three different columns can be co-located with at most one of them. There is no multi-column distribution and no second key to add. Any schema-wide plan is therefore an exercise in spending one scarce resource per table well, and in finding ways to make the other joins cheap by other means. ## Start with evidence The first mistake is designing from the ERD. Design from the workload: 1. Identify the queries that consume the most cluster time — not the most numerous, the most expensive in aggregate. 2. `EXPLAIN` them and read the distribution labels on each join. `DS_DIST_NONE` and `DS_DIST_ALL_NONE` mean nothing moved. `DS_DIST_INNER`, `DS_DIST_BOTH` and `DS_BCAST_INNER` each name a specific quantity of data crossing the network. 3. Estimate bytes moved per join per day: rows moved times row width times execution frequency. That number, not the number of queries, is what you are minimising. This usually collapses the decision quickly. Most schemas have one join that dominates the total by an order of magnitude. ## Spend the fact table's key Give the fact's DISTKEY to the join that moves the most bytes, subject to two vetoes: - **Cardinality veto.** The column must have many more distinct values than the cluster has slices, or most slices go unused. - **Uniformity veto.** The column must not have a dominant value, a sentinel like `-1`, or a large NULL population — all of which hash to a single slice and produce a straggler that slows every query on the table, not just the join you were optimising. If the biggest join's key fails a veto, take the next one. A distribution key that creates skew is worse than no co-location at all, because skew taxes every scan while a missing co-location only taxes the joins. ## Make the other joins cheap another way This is where most of the value actually is. `DISTSTYLE ALL` puts a full copy of a dimension on every node, so its joins are local no matter how the fact is distributed. For a dimension of a few hundred thousand rows this costs almost nothing in storage even on a wide cluster, and it converts the second and third join problems into non-problems. The budget question is storage times node count, plus slower writes since every load fans out to all nodes. That makes ALL right for small, slowly-changing dimensions and wrong for anything large or churning. A practical policy for a star schema: **every dimension is ALL unless it is large enough to hurt; the fact's single DISTKEY goes to the one large dimension that remains, or to the fact-to-fact join if one exists.** Fact-to-fact joins deserve special attention. Two large tables joined on a shared id — orders to order items, sessions to events — are where redistribution costs the most, and co-locating them on the shared key is often a better use of both tables' DISTKEYs than any dimension join. ## Where AUTO fits `DISTSTYLE AUTO` is the default and is managed by Redshift: small tables start as ALL, grow into EVEN, and Automatic Table Optimization can promote a table to KEY based on observed query patterns. Sort keys can be managed the same way. `SVV_ALTER_TABLE_RECOMMENDATIONS` exposes the recommendations the service has produced. The reasonable posture is a split policy. Leave exploratory, low-stakes and newly created tables on AUTO — it is better than an uninformed guess and it adapts. Pin explicit keys on the tables that dominate cost, because you know the join graph there and you do not want a background service changing the layout of your most expensive table without review. Read the recommendations either way; they are useful evidence even when you do not apply them. ## Sort keys are a separate budget Distribution decides whether the network is involved; the sort key decides how much of each slice's data has to be read. They are chosen for different reasons and should be reasoned about separately. In practice almost every large fact table wants a compound sort key led by its event timestamp, because analytical queries are time-bounded and appending in time order keeps maintenance cheap. Choosing the sort key to match the DISTKEY is only worth doing when you specifically want merge joins on a co-located key. ## Making it durable A key choice is a bet on a workload, and workloads move. Since `ALTER TABLE ... ALTER DISTSTYLE` and `ALTER SORTKEY` operate in place, treat the design as revisable: monitor `skew_rows` and `unsorted` in `SVV_TABLE_INFO`, re-read the plans of your top queries quarterly, and re-spend the keys when the shape of the workload has changed. The organisational half matters as much as the technical half — a review step on new large tables prevents most of the pathologies that later show up as a hot node.
- When would you spend a fact table's DISTKEY on another fact table rather than a dimension?When two large tables are joined regularly on a shared id — orders to order items, sessions to events. Redistributing two multi-terabyte tables is the most expensive movement in the cluster, while dimensions can usually be made local with DISTSTYLE ALL instead. Co-locating the large-to-large join is normally the higher-value use of both keys.
- How do you decide whether a dimension is too big for DISTSTYLE ALL?Compute its size times the node count, and weigh that against the redistribution it would otherwise cause. Also weigh write cost: every load fans out to all nodes, so a dimension rewritten frequently pays repeatedly. A few hundred thousand narrow rows is clearly fine; tens of gigabytes on a wide cluster clearly is not, and the middle needs measurement.
- Would you let Automatic Table Optimization manage your largest fact table?Usually not without review. AUTO is a good default for tables you cannot characterise, but on the table that dominates cluster cost you generally know the join graph better than an observer of recent queries, and you do not want the layout of your most expensive table changing unreviewed. Read the recommendations, apply them deliberately.
- How do you validate a distribution redesign before committing to it?Reproduce the top queries against the changed layout and compare both elapsed time and the plans' distribution labels, not just the skew ratio. A DISTKEY change usually trades one join's movement for another's, so the test has to cover the whole workload mix. Because ALTER DISTSTYLE works in place, rolling back is cheap if the measurement disappoints.
saying these in an interview costs you the question
- Trying to give a fact table more than one distribution key
- Designing keys from the ERD instead of from query evidence
- Replicating a large, frequently loaded dimension with DISTSTYLE ALL
- Choosing the biggest join's key even when that column is badly skewed
- Treating distribution and sort key as one combined decision