Two large tables in a sharded relational database must be joined, but rows with the same join value can live on different shards. What options do you have, and which do you pick for an OLTP request path?
answer
- Rows must meet on one machine to join
- Co-locate by one shared shard key
- Small read-mostly tables → replicate to every shard
- Denormalize: frozen copy vs cached copy
- Shuffle/broadcast joins = analytics, not OLTP
basics
~20 sOptions: co-locate both tables on the same shard key so the join is local; replicate small reference tables to every shard; denormalize the needed columns into the row; or join in the application after two lookups. For OLTP, co-locate or denormalize — never shuffle data between shards per request.
solid answer
~60 sRanked by what I would actually ship: 1. **Co-location** — shard both tables by the same key (for example both by `tenant_id`), so every row that can join lives on one shard and the join executes locally with normal engine plans. This is the design goal. 2. **Reference-table replication** — small, slow-changing lookup tables (currencies, plan types, country codes) get a full copy on every shard, so joins against them are always local. Cost: every write must fan out to all shards, so only do this for tables that are tiny and rarely written. 3. **Denormalization** — copy the two or three columns you actually needed into the row, so no join exists at request time. You now own keeping the copy fresh. 4. **Application-side join** — fetch from shard A, then fetch the matching rows from shard B by key. Fine for small, bounded row counts; it turns one query into N+1 round trips if you are careless. 5. **Broadcast or repartition (shuffle) join** — the engine moves data between nodes to align the join key. Legitimate for analytics, wrong for a user-facing OLTP path: unbounded latency and it makes the request depend on every shard.
code
sql · 23 lines-- both sharded by tenant_id, so the join is local on every shard
CREATE TABLE orders (
tenant_id BIGINT NOT NULL,
order_id BIGINT NOT NULL,
placed_at TIMESTAMPTZ NOT NULL,
PRIMARY KEY (tenant_id, order_id)
);
CREATE TABLE order_items (
tenant_id BIGINT NOT NULL, -- redundant, carried down on purpose
order_id BIGINT NOT NULL,
line_no INT NOT NULL,
sku TEXT NOT NULL,
unit_price_cents BIGINT NOT NULL, -- frozen copy: price at purchase time
PRIMARY KEY (tenant_id, order_id, line_no)
);
-- routable to one shard, local join
SELECT o.order_id, i.sku, i.unit_price_cents
FROM orders o
JOIN order_items i
ON i.tenant_id = o.tenant_id AND i.order_id = o.order_id
WHERE o.tenant_id = 42 AND o.order_id = 900123;go deeper
Say that the join only works cheaply if both rows sit on the same shard, and that the usual fixes are sharding both tables by the same key or copying the needed columns onto the row.
Rank the options and justify: co-location first, replicated reference tables for small read-mostly lookups, denormalization for hot reads, application-side join for bounded cases, shuffle joins only for analytics.
Talk about ownership of duplicated data — single writer, versioned idempotent propagation, reconciliation — and about which joins you deliberately push to a CDC-fed derived store to keep the OLTP fleet decoupled.
Frame the choice of the co-location dimension as the schema's defining decision, since it fixes which operations can stay transactional, and set policy for when duplication is allowed and who owns the consistency budget.
## Why the problem exists A join needs rows that share a value to meet in one place. In a single database they always do. Once a table is split across N servers by a shard key, two rows that join are on the same machine only if the sharding arrangement puts them there. If `orders` is sharded by `customer_id` and `products` by `product_id`, then `orders JOIN products` has its inputs scattered — no shard can answer alone. Relational engines have exactly two physical strategies for that situation, and both cost network: **broadcast** one side to every node, or **repartition (shuffle)** both sides on the join key. Distributed analytics systems do this routinely. In an OLTP request measured in single-digit milliseconds, either one is a disaster: latency becomes unpredictable, and the request now depends on every shard being healthy. So the practical answer is not "how do I execute a cross-shard join" but "how do I arrange the data so I never need one on the hot path." ## Option 1 — Co-location (the primary design tool) Pick one dimension that dominates access, and shard every table in that access path by it. In a business-to-business product that is usually the tenant or account; in consumer products it is often the user. `customers`, `orders`, `order_items`, `invoices`, `addresses` all sharded by `tenant_id` means any join among them is a local join on one shard, planned and executed by the ordinary engine with ordinary indexes. Most platforms make this explicit: you declare a group of tables that share a distribution key so the router can guarantee co-residency and, importantly, keep them together when data moves during rebalancing. Co-location is also what makes a multi-table *transaction* single-shard, which matters far more than the join itself. The limit: only one dimension can be co-located. A table joined heavily along two independent dimensions cannot be co-located for both. ## Option 2 — Reference (broadcast) tables Small, read-mostly tables that everything joins against — currency codes, product catalogue rows, feature flags, plan definitions — get a **full copy on every shard**. Joins against them are then always local, whatever the other table's shard key is. The trade is on writes: every insert or update must reach all N shards, which means either a fan-out write, a periodic full refresh, or a replication feed. That is acceptable for a table that changes daily and holds thousands of rows; it is a trap for anything that grows or is written per request. Watch for a reference table that quietly becomes large — a "catalogue" that reaches millions of rows costs storage on every shard and rebuilds on every reshard. ## Option 3 — Denormalization Instead of joining to fetch `product.name` and `product.price_cents`, store those values on the `order_item` row at write time. The join disappears from the read path entirely. This is the standard sharded-OLTP move, and it has two flavours worth distinguishing: - **Point-in-time copies** that are *supposed* to be frozen — the price at the moment of purchase. These are not stale data; they are the correct value, and the source table changing later must not change them. - **Cached copies** of mutable attributes — a display name duplicated for convenience. These need an update path: one owner writes the source, propagates asynchronously (change feed, outbox), applies idempotently with a version so late messages cannot overwrite newer values, and a reconciliation job that repairs drift. Confusing the two is the classic bug: someone "fixes" the frozen historical price by backfilling it from the current catalogue. ## Option 4 — Application-side join Query shard A, collect the foreign keys, query shard B (or several shards) by those keys, stitch in memory. Correct and simple, and if the second lookup is a batched `IN (...)` against a bounded key set it is fast. The failure mode is the N+1 pattern — one round trip per parent row — and unbounded key sets that quietly become a fan-out anyway. Batch the second call, cap the key count, and set timeouts. ## Option 5 — Let the engine shuffle Repartition and broadcast joins belong to analytics: a reporting replica, a columnar store, or a warehouse fed by change data capture. Route the "join everything by everything" questions there rather than teaching the OLTP shards to do it. That also protects the transactional fleet from a single analyst's query pinning all N nodes. ## How to choose Ask: is this join on a user-facing hot path? If yes, it must be local — co-locate, replicate a small reference table, or denormalize. If it is rare and back-office, an application-side join or a bounded fan-out is fine. If it is analytical, it does not belong on the shards at all.
- You denormalized a product name onto each order line. Marketing renames the product. What has to happen, and what must not?Decide first whether the copy is a frozen point-in-time fact or a cached attribute. If order lines are meant to show what the customer bought at the time, the old name stays and nothing propagates. If it is a cache of the current name, the product service publishes a change event carrying a version, each shard applies it idempotently and only when the version is newer than what it holds, and a periodic reconciliation job compares copies against the source to repair anything dropped.
- Why is co-location worth more than just making joins local?Co-located tables also let a multi-table write stay inside one shard, which means one local ACID transaction instead of a distributed commit. Avoiding cross-shard transactions is usually the larger win, because the alternatives — two-phase commit or sagas with compensation — cost far more in latency, availability, and code complexity than a slow join would.
saying these in an interview costs you the question
- Proposing that the router 'just does the join' without accounting for the data movement between shards
- Replicating a large or frequently written table to every shard as a reference table
- Treating every denormalized column as a cache that must be kept current, including values meant to be frozen at write time
- Doing an application-side join one parent row at a time (N+1) instead of a batched key lookup
- Assuming a second shard key can be added so a table is co-located along two dimensions at once