When designing a sharded transactional schema, how do you decide between duplicating data so every request stays on one shard versus accepting cross-shard fan-out for the secondary access pattern?
answer
- Only one dimension can be co-located
- Fan-out cost grows with fleet size; duplication cost grows with change
- Fan-out re-couples independent failure domains
- Directory hop < payload copy when latency allows
- Every duplicate is a future backfill
basics
~20 sDecide per access path by rate and latency target. Hot user-facing paths must stay single-shard, so duplicate the data and pay write amplification and a consistency window. Rare back-office paths can fan out. Anything analytical belongs in a derived store, not on the shards.
solid answer
~60 sI budget it per access path rather than deciding globally. **Inputs:** queries per second, latency objective, tolerance for staleness, and the write rate on the duplicated attribute. **Fan-out costs** compound with fleet size: latency is the slowest shard's, capacity for those queries never grows with more shards, and the request depends on *all* shards, so a single sick node degrades every fan-out path. That last point is architectural — fan-out converts independent failure domains back into one. **Duplication costs** are steady-state and mostly bounded: write amplification, a staleness window, storage, and — the one people forget — every schema change to a duplicated attribute becomes a backfill across the fleet. **My rules:** the dominant access dimension picks the shard key. The second dimension gets either a directory table turning it into two point lookups, or a full duplicate keyed by it with one writer, versioned idempotent propagation, and a reconciler. Reporting and search go to a change-data-capture-fed store. Whatever fan-out remains is capped — bounded concurrency, per-shard timeouts, admin-only — and tracked as a fan-out budget with an explicit owner.
go deeper
Know the basic trade: copying data avoids querying every shard but means the copy can go stale and must be kept up to date.
Compare the concrete costs — fan-out latency and lost throughput scaling versus write amplification, staleness, and storage — and say that hot paths should stay single-shard.
Decide per access path with numbers, specify the consistency contract for each duplicate (single writer, versioned idempotent propagation, stated staleness, reconciliation), and move analytical queries to a change-data-capture-fed store.
Lead with the availability argument — fan-out re-couples failure domains and its cost grows with every shard added — set an enforced fan-out budget with owners, and distinguish the irreversible co-location decision from the revisable duplication ones.
## Why this is a judgment call A sharded schema can be co-located along exactly one dimension. Real products have at least two: orders by customer *and* by warehouse; messages by conversation *and* by author; tickets by tenant *and* by assignee. One of those is local and cheap; the other is either fan-out or a copy. There is no arrangement that makes both free, so the work is deciding which cost to pay where. ## The two cost curves **Fan-out costs scale badly with fleet size and are paid at request time.** - Latency is the maximum over N shards, so per-shard tails become per-request norms. A 1% slow rate per shard is a 27% slow rate across 32. - Capacity does not scale: every shard sees 100% of fan-out queries, so growing the fleet adds no headroom for them while adding load. - Availability compounds. A request touching all shards has availability roughly a^N. This is the deepest cost, because sharding's value proposition is *independent failure domains*, and a hot fan-out path silently re-couples them. After a reshard the exposure gets worse, not better. - Blast radius: one degraded shard becomes a whole-product incident. **Duplication costs are steady-state and mostly bounded.** - Write amplification: one logical write becomes two or more. - A staleness window between the source and the copy, whose size is a product decision. - Storage and index cost, usually the cheapest term. - Operational drag: propagation pipelines to run, drift to reconcile, and — the term teams consistently underestimate — every future change to a duplicated attribute becomes a fleet-wide backfill. A copy is a permanent migration liability, not a one-time write. - Correctness surface: duplicated data can disagree, and someone will read the wrong copy. The asymmetry is that fan-out costs grow with success (more traffic, more shards) while duplication costs grow with change (more schema evolution). Fast-growing systems should therefore lean toward duplication; stable systems with rare secondary access can tolerate fan-out. ## The decision procedure **1. Enumerate access paths, not tables.** For each: rate, latency objective, staleness tolerance, who calls it. **2. Pick the co-location dimension from the highest-rate latency-sensitive cluster of paths.** Everything transactional that mutates together must sit inside it, because single-shard transactions matter more than single-shard reads. **3. Classify each remaining path:** - *Hot and user-facing* — must be single-shard. Choose the cheapest mechanism that achieves it: a **directory table** mapping the secondary key to the shard key (two point lookups, no duplicated payload, no staleness on the payload itself) if that extra hop fits the budget; otherwise a **full duplicate row** keyed by the second dimension. - *Warm and internal* — a bounded fan-out with concurrency caps and per-shard timeouts is acceptable if the merged result is small. - *Analytical or search-shaped* — do not put it on the shards at all. Feed a search index, a columnar store, or a warehouse from change data capture. This also protects the transactional fleet from ad-hoc queries. **4. Prefer the directory over the payload copy** when you can afford the hop. It duplicates an immutable mapping rather than mutable attributes, so there is far less to keep consistent and far less to backfill later. **5. Set the consistency contract explicitly.** For every duplicate: one writer owns the source; propagation is asynchronous, idempotent, and version-guarded so late messages never overwrite newer values; the tolerated staleness is a stated number; a reconciliation job compares copies and repairs or alerts; and it is documented which copy is authoritative when they disagree. **6. Budget and enforce the remainder.** Track the share of QPS that touches more than one shard, per path. Make new fan-out paths a design-review decision rather than something that appears in a pull request. Cap degree of fan-out, set timeouts, and consider whether partial results are acceptable — they usually are for search-like features and never for money. ## Second-order considerations - **Reshard cost.** Every duplicate is more data to move and more invariants to verify when buckets migrate. Duplication makes the fleet cheaper to run and more expensive to reshape. - **Tenant skew.** In a multi-tenant fleet, a single large tenant can make a "cheap" fan-out expensive on one shard, or make a duplicate hot. Design for the p99 tenant, not the mean. - **Team topology.** A duplicate crossing service boundaries needs a clear owner for the propagation pipeline; unowned pipelines rot into drift. - **Reversibility.** Adding a duplicate later is a backfill — doable. Changing the co-location dimension later is a full reshard of the schema — usually not. Spend disproportionate care on the dimension choice, and treat duplication decisions as revisable. ## How to state the answer "Single-shard for the hot path, always. Second dimension via a directory hop if latency allows, a duplicate if it does not. Everything analytical leaves the OLTP fleet through change data capture. Remaining fan-out is capped, owned, and budgeted — and I care most about the availability coupling, because that is the cost that grows every time we add a shard."
- Why prefer a directory table over duplicating the whole row, when both remove the fan-out?A directory duplicates only a mapping from the secondary key to the shard key, which is typically immutable, so there is almost nothing to keep consistent and nothing to backfill when the row's columns change. A full payload copy duplicates mutable attributes, so it needs a propagation pipeline, versioned idempotent application, a staleness contract, and a reconciler, and it turns every future schema change into a fleet-wide migration. The directory's cost is one extra round trip, so you choose it whenever the latency budget can absorb that hop.
- What signal tells you a fan-out path has become an architectural problem rather than an acceptable cost?Watch its share of total QPS and its correlation with incidents. When a single degraded shard measurably moves the error rate or p99 of a user-facing endpoint, the path has re-coupled your failure domains and must be removed. The other signal is trend: fan-out cost rises every time you add a shard, so a path that is merely uncomfortable at eight shards is an outage generator at thirty-two.
- How does duplication interact with the cost of resharding later?Every duplicate adds data to copy and invariants to verify during a bucket move, and the propagation pipeline itself must keep working correctly while ownership changes hands mid-flight. So duplication trades cheaper steady-state operation for a more expensive and more delicate reshape. That is usually the right trade, but it argues for keeping the number of distinct duplicated projections small and each one clearly owned.
saying these in an interview costs you the question
- Deciding duplication versus fan-out once for the whole system instead of per access path
- Treating duplication as free because storage is cheap, ignoring write amplification and future backfills
- Ignoring that a fan-out request's availability is the product of all shards' availability
- Putting reporting or search queries on the transactional shards because 'the data is already there'
- Assuming the co-location dimension can be changed later as easily as a duplicate can be added