Hiring BarSupport

Design a Sharded Database — Staff-Level Case Study

Case study66 min read8 diagrams

Technologies referenced in this case study: PostgreSQL · DynamoDB · Cassandra · ZooKeeper & etcd · Apache Kafka

Related: Sharding & Partitioning · Consistent Hashing · Scaling Writes · Database Selection · Replicated Data Store

How to Use This Case Study#

Organized for interview use first, reference second. Read front-to-back once. Return to individual sections for targeted review.

ModeTimeWhat to Read
Quick Review15 minExecutive Summary → Interview Walkthrough → Fault Lines table → Active Drills 1–3
Targeted Study1–2 hrsExecutive Summary → Walkthrough → Section 3 (Fault Lines) → Section 4 (Failures) → weak-spot Deep Dives
Deep Dive3+ hrsEverything, including the Principal Lens and appendices
What is Database Sharding? — Why interviewers pick this topic

Sharding splits one logical dataset across many independent database nodes, each owning a disjoint subset of rows chosen by a shard key. Each shard is a full database with its own primary, replicas, WAL and backups. The application (or a routing layer) decides which shard serves each query.

Before vs After — Primary-at-the-ceiling scenario:

Without sharding (single PostgreSQL primary, 64 vCPU, 4 TB):
t=0:       Write QPS crosses 18K sustained; WAL generation at 250 MB/s
t=+2wk:    Replica lag spikes to 40s during nightly batch; reads go stale
t=+1mo:    Autovacuum can no longer keep up on the orders table (1.8B rows)
t=+6wk:    Largest instance class already in use — nowhere to scale up
t=+2mo:    Black Friday: primary CPU 100%, p99 write latency 1.2s, checkout errors
t=+2mo+4h: Incident closes with "we need to shard" — now as an emergency project

With sharding (planned 9 months earlier, 32 logical shards on 8 hosts):
t=0:       Same 18K write QPS spread ~2.3K per host
t=+2mo:    Black Friday: hottest host at 55% CPU; p99 write 14ms
t=+2mo:    Shard 17 runs hot (large merchant) — moved to its own host in 40 min
t=+1yr:    Hosts doubled to 16 by moving logical shards — zero application changes

Why interviewers reach for this question: Sharding is the first design decision most engineers make that is genuinely a one-way door. The shard key choice outlives the team that picked it. Interviewers use it to see whether you can reason about access patterns, data movement, cross-shard correctness, and the multi-year operational tax — not whether you can draw a hash ring.

Mechanics Refresher: Partitioning Strategies
StrategyHow It WorksProsCons
Hash (mod N)shard = hash(key) % NEven distribution, trivial routingChanging N remaps ~100% of keys
Consistent hashingKeys and nodes on a ring; key goes to next node clockwiseAdding a node moves ~1/N of keysUneven without virtual nodes; range scans impossible
RangeContiguous key ranges per shard ([a, m), [m, z))Efficient range scans; easy splitsHot tail on monotonic keys (timestamps, auto-increment IDs)
Directory / lookupA mapping table key → shardArbitrary placement; move individual tenantsThe directory is a new critical dependency
Logical shards → physical hostsHash into fixed M logical shards (e.g., 4,096), map logical → physicalRebalance by moving whole logical shards; N changes without rehashingM is itself a one-way door; too small caps scale
Geo / tenantShard by region or tenant IDData residency, tenant isolationWildly uneven shard sizes

For most production systems: hash the shard key into a fixed number of logical shards (1,024–4,096), and map logical shards to physical hosts via a small, versioned directory. You get hash-level evenness with directory-level flexibility. The strategy is almost never the interview question — the shard key, resharding, and cross-shard correctness are.


Executive Summary

If you only read one section, read this. Everything in the case study flows from the contrast below.

What This Interview Actually Tests#

Sharding is not a partitioning-algorithm question. Everyone can draw consistent hashing.

It is an access-pattern and data-movement ownership question that tests:

  • Whether you derive the shard key from the top 3 queries, not from the schema
  • Whether you know which queries you are about to make expensive (or impossible)
  • Whether you have a plan for moving data while it is being written
  • Whether you can name who owns the router, the directory, the resharding runbook and the cross-shard invariants

The key insight: The shard key decides which queries stay single-shard for the next five years. Staff engineers pick it by listing the queries that must never fan out, then design resharding before they need it — because the first resharding is always done under duress.

The L5 vs L6 Contrast — Start Here#

BehaviorSenior (L5)Staff (L6)Principal (L7)
First move"We'll shard by user_id with consistent hashing""What are the top 3 queries, and which must stay single-shard?""Should we shard at all, or buy a distributed SQL engine? What's the 3-year TCO of each?"
Shard keyPicks the primary keyPicks the key that co-locates the transactional boundary; names the queries that will fan outStandardizes a tenant/entity-scoped key convention across 20 services so cross-service joins stay aligned
Resharding"Consistent hashing handles it"Logical shards + online migration (dual-write/backfill/verify/cutover) with a rollback at every stepMakes resharding a platform capability with an SLO (e.g., "move a 200 GB shard in < 2h, zero write loss")
Cross-shard"Use two-phase commit"Designs so 95%+ of transactions are single-shard; sagas/outbox for the rest; names the invariant that can driftDefines org-wide policy: which invariants may be eventually consistent, who signs off, how drift is audited
Hot shards"Add more shards"Distinguishes hot shard vs hot key; isolates whales onto dedicated shards via directory overridePrices tenant isolation into the contract; enterprise tier gets a dedicated cell
OwnershipApplication team owns routing logicRouting in a shared library or proxy (Vitess-style); DB platform owns the directoryRedraws boundaries: data platform owns shard topology; product teams own keys and schemas only
Why "first move" separates levels

L5: Leads with mechanism: "shard by user_id, consistent hashing, 16 shards." It is often the right answer — for a consumer app. But it is chosen before anyone has said what the system does, so the candidate cannot defend it when the interviewer reveals the main query is "all orders for a merchant in the last 30 days."

L6: "Before picking a key I want the three queries that pay the bills. If the hot path is 'load a user's cart and write an order,' the transactional boundary is the user, and user_id is the key. If it's 'merchant dashboard over their orders,' merchant_id is the key and the user view becomes a secondary index. I'll assume a marketplace where checkout is the hot path."

L7: Asks whether sharding is the right investment at all. A 4 TB PostgreSQL at 15K writes/s can often be bought another 2 years with vertical partitioning, archival, and a read-replica tier — or replaced by a managed distributed SQL engine. The L7 candidate prices the options.

Why "resharding" separates levels

L5: "Consistent hashing means adding a node only moves 1/N of the keys." True, and irrelevant: moving 1/N of 20 TB is 1.25 TB of live, mutating data. The hard part is not which keys move — it's moving them without losing writes, without doubling latency, and with a rollback.

L6: Separates the placement problem (logical → physical map) from the data-movement problem (snapshot + change stream + verification + atomic cutover), and states the cutover's write-unavailability budget out loud: "Cutover is a 1–5 second write pause per logical shard, done during low traffic, with a single config flip back if verification fails."

Why "cross-shard" separates levels

L5: Reaches for 2PC because it is the textbook answer to "transaction across two databases." 2PC works, but it holds locks across a network round-trip, blocks on coordinator failure, and couples the availability of every shard in the transaction.

L6: Treats cross-shard transactions as a data modeling failure to be minimized, then handles the residue explicitly: "I'll pick the key so 98% of writes are single-shard. For the 2% — transfers between two users — I'll use an outbox + saga with idempotent steps, and I'll say which invariant is temporarily violated: the sum of balances can be off by in-flight transfers for up to a few seconds."

The Staff Positions#

PositionRationale
Derive the shard key from the transactional boundaryThe entity that must be updated atomically (user, tenant, order) should never span shards
Logical shards ≫ physical hosts (1,024–4,096 logical)Rebalancing becomes "move a shard," never "rehash the world"
Don't shard until you must — then shard early enoughSharding before ~5–10 TB or ~20K writes/s buys operational pain for no gain; sharding at 95% CPU is an emergency
Avoid distributed transactions by design, not by protocol2PC is a last resort; sagas + outbox for the residue; name the drifting invariant
Directory override for whalesHash for the long tail, explicit placement for the top 0.1% of tenants
Resharding is a product with an SLOOnline migration tooling built and exercised before the first emergency
Global secondary indexes are asyncMaintained via CDC, not synchronous cross-shard writes; readers accept seconds of lag

The Three Intents#

Three intents drive every design decision. Each leads to a different architecture.

IntentConstraintStrategyFailure ModeCorrectness Bar
Write-throughput scaling (OLTP)One primary can't absorb write QPS or dataset sizeHash on entity key into logical shards; single-shard transactionsHot shards; cross-shard transactionsStrong consistency within a shard; eventual across
Multi-tenant isolation (SaaS)Noisy tenants hurt others; enterprise needs isolation/residencyTenant-ID key, directory placement, dedicated shards for whalesUneven shard sizes; tenant outgrows a shardTenant data never leaks; per-tenant SLOs
Time-series / append-heavyIngest volume; old data rarely readRange-by-time + hash-within-bucket; TTL by dropping partitionsHot "now" partition; cross-time queries fan outCompleteness per window; retention compliance

🎯 Staff Move: "I'll assume OLTP write scaling for a marketplace, because that's where shard-key choice and cross-shard transactions collide. If this were multi-tenant SaaS I'd key on tenant_id with a directory; if it were time-series I'd key on time buckets — and those are three different designs with different failure modes."

The Five Fault Lines#

#Fault LineThe Tension
1Shard Key: Locality vs DistributionCo-locate related rows for single-shard queries, or spread evenly to avoid hot shards?
2Placement: Algorithmic vs DirectoryCompute location (no dependency, rigid) or look it up (flexible, new critical path)?
3Resharding: Planned Online Migration vs Emergency SplitInvest in migration tooling up front, or pay for it during an incident?
4Cross-Shard Correctness: Distributed Transactions vs SagasAtomicity with coupled availability, or availability with temporary invariant drift?
5Routing Ownership: Library vs Proxy vs EngineApp-embedded routing, a sharding proxy (Vitess-style), or buy a distributed SQL engine?

In the Wild: Real Production Systems#

Why this section belongs here: Citing specific production systems demonstrates you've studied operational reality, not textbook designs.

YouTube / Vitess — Sharding Proxy over MySQL#

Vitess was built at YouTube to scale MySQL horizontally and is now a CNCF project used by Slack, GitHub and others. Applications talk to VTGate, a stateless proxy that parses SQL and routes by a vindex (the sharding function) defined in a VSchema. Resharding is performed online by the VReplication-based workflows (Reshard, MoveTables): copy, catch up via binlog streaming, verify with VDiff, then switch reads and writes.

Staff insight: Vitess productizes exactly the steps an interviewer wants to hear — copy, catch-up, verify, cutover, reverse replication for rollback. Name those phases and you've shown you understand resharding as a workflow, not an algorithm.

Instagram — Logical Shards in PostgreSQL Schemas#

Instagram's well-known engineering post describes mapping thousands of logical shards to a much smaller number of physical PostgreSQL servers, with each logical shard as a PostgreSQL schema. IDs are generated in-database and embed the logical shard ID alongside a timestamp, so any ID can be routed without a lookup.

Staff insight: Encoding the logical shard in the ID makes routing stateless and lets you move logical shards between hosts without rewriting IDs. That's the "logical ≫ physical" position in production.

Notion — Sharding PostgreSQL by Workspace#

Notion publicly described sharding its PostgreSQL fleet into 480 logical shards across 32 physical databases, partitioned by workspace ID so a workspace's blocks stay together, and later re-sharding onto more hosts using a dual-write/backfill/verify approach.

Staff insight: Workspace is the transactional and access boundary, so it's the shard key. 480 was chosen because it has many divisors — letting the logical shards spread evenly across many physical host counts. Numbers like that are what an L6 answer sounds like.

What Interviewers Probe#

After You Say...They Will Ask...(What They're Evaluating)
"Shard by user_id""How does a merchant see all their orders?"Do you know which queries you just made fan-out?
"Consistent hashing""You're adding 4 hosts to 8. Walk me through the next 6 hours."Can you move live data without losing writes?
"Two-phase commit""Coordinator dies after prepare. What's locked, and for how long?"Do you understand the availability cost of atomicity?
"We'll add more shards""One tenant is 30% of traffic. Does that help?"Hot shard vs hot key distinction
"Global secondary index""Is it updated synchronously? What happens on partial failure?"Index consistency and write amplification
"The router knows the mapping""Where is the mapping stored, and what happens when it's stale?"Control-plane failure modes

System Architecture Overview#

Diagram: System Architecture Overview

Reading the diagram: The router is on every request but holds no durable state — it reads a cached, versioned shard map. The directory in etcd is on the control plane and can be down for minutes without taking out the data plane, as long as the cached map is valid. Everything that is not keyed by the shard key — lookups by buyer, analytics, search — is fed asynchronously by CDC and is explicitly allowed to lag.

Quick-Reference: The 30-Second Cheat Sheet#

TopicThe L5 AnswerThe L6 Answer — Say This
When to shard"When the DB gets big""When write QPS or working set exceeds one primary with 2× headroom — roughly 20–50K writes/s or 5–10 TB on commodity hardware. Before that: archive, partition, replicate."
Shard key"user_id""The entity that forms the transactional boundary for the top queries. I'll list which queries fan out as a result."
Placement"Consistent hashing""Hash into 4,096 logical shards, map logical → physical in a versioned directory. Rebalance by moving logical shards."
Resharding"Add nodes, ring rebalances""Snapshot + CDC catch-up + VDiff-style verify + 1–5 s write-freeze cutover, with reverse replication for rollback."
Cross-shard writes"2PC""Design them out; outbox + saga for the residue; name which invariant drifts and for how long."
Hot shard"Add shards""Adding shards doesn't split a key. Move the whale to a dedicated shard via directory override; split its key space only if it's truly one hot row."

Key Numbers Worth Memorizing#

MetricValueWhy It Matters
Comfortable single PostgreSQL/MySQL primary~10–30K writes/s, 2–10 TBBelow this, don't shard — tune, archive, partition
Target size per physical shard200 GB – 1 TBRestore/rebuild time stays under ~1–2 hours
Logical shards1,024–4,096 (or highly composite, e.g., 480, 720)Enough granularity to rebalance for 5+ years
Router overhead0.3–1 ms per queryThe proxy tax; cheap compared with fan-out
Fan-out query latencyp99 ≈ max of N shard p99s; with 64 shards, tail dominatesWhy scatter-gather must be rare
Probability at least one of 64 shards is slow (each 1% slow)1 − 0.99⁶⁴ ≈ 47%Fan-out turns rare slowness into common slowness
Online move throughput~50–200 MB/s per streamA 500 GB shard copies in ~45–170 min
Cutover write-freeze1–5 s per logical shardBudget you state to product up front
2PC added latency+1–2 RTTs, locks held across them5–20 ms intra-region; unbounded on coordinator failure
Consistent hashing move on +1 node~1/(N+1) of keysBut still TBs of live data at scale
Hash mod N remap on N → N+1~N/(N+1) of keys (≈ all)Why naive modulo is a trap

Interview Walkthrough

The most common mistake: Candidates spend 20 minutes on hash functions, ring diagrams and replica counts, then get asked "how do you add capacity?" with 8 minutes left. Compress the basics to ~10 minutes; the level is decided in the shard-key defense, the resharding workflow and cross-shard correctness.


Phase 1: Requirements & Framing (2–3 min)#

State the functional scope in one sentence:

"We're storing orders for a marketplace — create, update status, read by buyer, read by merchant — and the single primary is running out of write headroom."

Then spend your time on the non-functional constraints that pick the design:

"Three numbers matter: write QPS and growth rate, dataset size and growth, and which queries must stay single-shard with strong consistency. I'll assume 40K writes/s peak growing 2× a year, 12 TB today, and that checkout — creating an order and decrementing a merchant's inventory — must be transactional."

Commit to an intent:

"I'll treat this as OLTP write scaling. Multi-tenant isolation and time-series retention are different designs; I'll point out where they'd diverge."

🎯 Staff Move: Ask for the top three queries by business value, not by volume. The shard key serves the queries that must be fast and correct; everything else can be served by an async index or the warehouse.


Phase 2: Core Entities & API (1–2 min)#

  • Merchant (merchant_id) — owns inventory and orders; the transactional boundary for checkout
  • Order (order_id, merchant_id, buyer_id, status, total, created_at)
  • InventoryItem (merchant_id, sku, qty)
  • ShardMap (logical_shard → physical_cluster, overrides: merchant_id → cluster, version)
CreateOrder(merchant_id, buyer_id, items[])     → order_id        # single-shard txn
GetOrder(order_id)                              → Order           # routable via embedded shard id
ListOrdersByMerchant(merchant_id, since, limit) → Order[]         # single-shard
ListOrdersByBuyer(buyer_id, limit)              → Order[]         # async GSI, may lag ~1–2 s

"Order IDs embed the logical shard — order_id = [41 bits time][12 bits logical shard][10 bits sequence] — so GetOrder never needs a lookup."


Phase 3: High-Level Architecture (≤5 min)#

Diagram: Phase 3: High-Level Architecture (≤5 min)

Walk the flow in 60 seconds:

  1. Every query carries merchant_id (or an order_id that encodes the logical shard).
  2. The router hashes to one of 4,096 logical shards and looks up the physical cluster in a cached map.
  3. Checkout — insert order + decrement inventory — is a local transaction on one shard.
  4. Each shard's WAL/binlog is streamed via CDC into Kafka, which feeds the buyer-lookup index and the warehouse.
  5. The directory is control plane: routers cache the map and keep serving if etcd is briefly unavailable.

🎯 Staff Move: "This is the design that works on day one. What decides whether it works in year three is how we reshard, what happens to the queries I just made expensive, and how we handle the merchant who becomes 20% of traffic." You're at ~10 minutes.


Phase 4: Transition to Depth (1 min)#

"Three areas worth going deep on: defending the shard key against the queries it hurts, how we reshard online without losing writes, and cross-shard correctness for the few operations that span merchants. I'd start with resharding — it's the one-way door that most designs get wrong. Which would you like?"


Phase 5: Deep Dives (25–30 min)#

Deep dive 1: Shard key defense (5–7 min)

Enumerate what the key costs:

QuerySingle-shard?How it's servedCost
Checkout (order + inventory)✅Local transactionNone
Merchant dashboard✅Index on (merchant_id, created_at)None
Get order by ID✅Logical shard embedded in IDNone
Buyer's order history❌Async GSI keyed by buyer_id1–2 s lag; +1 write per order
Platform-wide GMV report❌Warehouse via CDCMinutes of lag
Admin search by email❌Search index via CDCSeconds of lag

"Four of six hot queries stay single-shard. The two that don't tolerate seconds of lag. That's why merchant_id wins over buyer_id: keying by buyer would make checkout cross-shard, because inventory belongs to the merchant."

Deep dive 2: Online resharding (8–10 min) — covered in 3.3. Name the phases: copy → catch-up → verify → freeze → flip → reverse-replicate → cleanup. State the write-freeze budget (1–5 s) and the rollback (flip the map back; reverse replication has kept the old shard current).

Deep dive 3: Cross-shard correctness (5–7 min) — a buyer uses store credit issued by the platform (different shard). "Outbox on the order shard, saga step debits credit idempotently keyed by order_id, compensation on failure. The invariant 'credit ledger balances' is eventually consistent for up to ~5 s; finance signs off because reconciliation runs hourly."

Deep dive 4: Hot shards (3–5 min) — directory override for whales; split a logical shard only if it contains multiple heavy keys; true single hot rows need write sharding (counter splitting) in the schema.


Phase 6: Wrap-Up (2–3 min)#

"To summarize: merchant_id into 4,096 logical shards over 16 clusters, directory-backed placement with whale overrides, async GSIs for buyer queries, online resharding as a tested workflow, and sagas for the ~2% of cross-merchant operations. What I'd build next: automated shard-balancing based on shard.cpu and size, and a quarterly resharding game day. What I would not build: global secondary indexes with synchronous writes, or a general-purpose cross-shard query engine — that's what the warehouse is for."

Common Timing Mistakes#

MistakeTime LostFix
Explaining consistent hashing from first principles5–8 min"Hash into logical shards; the ring is an implementation detail."
Designing replication inside each shard4–6 min"Each shard is a standard primary + 2 replicas. Same as unsharded."
Building a distributed query planner5+ min"Cross-shard analytics go to the warehouse."
Never stating which queries fan out—The interviewer will ask; say it first
Leaving resharding for the last 2 minutes—It's the highest-signal topic; pull it forward

1. The Staff Lens#

1.1 Why This Problem Exists in Staff Interviews#

Sharding is where local design decisions become organizational constraints. Once 20 services join on merchant_id and 3 years of IDs encode logical shard numbers, changing the key is a multi-quarter migration. Interviewers want to see that you know you're choosing a constraint for teams who haven't been hired yet — and that you design the escape hatch (logical shards, directory overrides, tested online migration) before you need it.

1.2 The L5 vs L6 Contrast — Visual#

Diagram: 1.2 The L5 vs L6 Contrast — Visual

1.3 The Staff Question That Cuts Through Everything#

"Which entity must be updated atomically, and which queries are we willing to make eventually consistent?"

Every sharding decision follows from the answer. The transactional boundary becomes the shard key. The queries on other dimensions become async indexes with an explicit lag budget. If the answer is "everything must be atomic with everything," you're not designing a sharded database — you're buying a distributed SQL engine and paying its latency.


2. Problem Framing & Intent#

2.1 The Three Intents — Explained#

Write-throughput scaling (OLTP). The primary is the ceiling. The key is the entity at the center of the write path (user, merchant, account). Correctness bar: ACID within the entity. Everything else — search, reporting, cross-entity views — moves to derived stores. Typical scale: 20K–500K writes/s across 16–256 clusters.

Multi-tenant isolation (SaaS). The key is tenant_id with a directory, because tenant sizes follow a power law: in a typical B2B SaaS, the top 1% of tenants hold 30–50% of data. Hash placement for the long tail; explicit placement (dedicated clusters, specific regions for residency) for the head. The failure mode is a single tenant outgrowing a shard — which then requires intra-tenant sharding, a second key.

Time-series / append-heavy. Range on time buckets (daily/hourly) with hash within the bucket. Retention becomes DROP PARTITION, not DELETE — orders of magnitude cheaper. The failure mode is the hot "now" bucket; the fix is hashing within it ((day, hash(device_id) % 64)).

🎯 Staff Move: "These three have different keys, different placement and different failure modes. If the product has both OLTP orders and an event log, I'd shard them separately — one sharding scheme per access pattern, not one per company."

2.2 When NOT to Use Database Sharding#

SituationBetter AlternativeWhy
Read-heavy load, writes < 10K/sRead replicas + cache (Scaling Reads)Sharding solves write/size limits, not read limits
Dataset < 2–5 TB, growing slowlyNative table partitioning + archivalSame pruning benefits, one database, no router
Need arbitrary multi-row transactions across entitiesDistributed SQL (Spanner, CockroachDB, Yugabyte)Pay latency for automatic sharding + distributed txns
Key-value access at huge scaleDynamoDB / CassandraSharding is built in; don't rebuild it on MySQL
One huge table dominates, rest is smallVertical split: move that table to its own storeShard one table, not the whole schema
Team of 4 with no DB platformManaged service with auto-partitioningHome-grown sharding needs a platform owner

"The cheapest shard is the one you don't create. I'd first move the 60% of data that's cold order history to object storage, and see if that buys two years."

2.3 What the Interviewer Leaves Underspecified#

GapWhy It MattersWhat to Say
Growth rateDecides logical shard count"At 2×/year, 4,096 logical shards covers ~8 doublings from 16 hosts."
Transaction scopeDecides the key"Is checkout atomic across merchants (a multi-merchant cart)?"
Secondary access patternsDecides GSIs vs fan-out"How often does support look up by email? Seconds of lag OK?"
Tenant size distributionDecides hash vs directory"Is there a Shopify-sized merchant, or is it long tail?"
Data residencyForces geo placement"Any EU data that must stay in the EU?"
Existing systemMigration vs greenfield"Are we migrating a live monolith? That's the real project."

2.4 Precise Terminology#

TermPrecise MeaningCommon Confusion
PartitioningSplitting a table's rows (can be within one database)Used interchangeably with sharding
ShardingPartitioning across independent database nodesNot replication — shards hold different data
Shard keyColumn(s) whose value determines placementNot necessarily the primary key
Logical shardUnit of placement and movement; fixed countConfused with a physical host
Directory / shard mapVersioned mapping logical → physical (+ overrides)Not a per-row lookup table
Scatter-gather / fan-outQuery sent to all (or many) shards, results mergedAssumed "just slower" — actually a tail-latency amplifier
Local secondary indexIndex within each shard, same partitioningRequires fan-out when queried without the shard key
Global secondary index (GSI)Index partitioned by the indexed columnRequires cross-shard writes or async maintenance
CutoverAtomic switch of writes from old to new locationNot the same as "copy finished"
Hot shard vs hot keyMany busy keys on one shard vs one key that is busyAdding shards fixes the former, not the latter

3. The Five Fault Lines#

Each fault line: the options, who pays, the Staff default, and when to deviate.

3.1 Fault Line 1: Shard Key — Locality vs Distribution#

The tension: A key that co-locates related data (all of a merchant's orders) makes transactions and range queries single-shard but concentrates load. A key that spreads evenly (hash of order_id) balances load but makes every entity-scoped query a fan-out.

Key ChoiceWhat WorksWhat BreaksWho Pays
order_id (random)Perfect balanceMerchant dashboard = 4,096-way fan-out; checkout txn cross-shardProduct teams (every entity query slow)
merchant_idCheckout + dashboard single-shardWhale merchants create hot shards; buyer view fans outDB platform (whale placement) + buyer-view owners (lag)
buyer_idBuyer history single-shardInventory decrement becomes cross-shardCheckout team (sagas)
(merchant_id, order_month) compoundSpreads whale over timeMerchant txns spanning months cross-shard; range queries fan out by monthEveryone reading >1 month
created_at rangeTime-range scans100% of writes on the newest shardOn-call (hot tail)
Diagram: 3.1 Fault Line 1: Shard Key — Locality vs Distribution

Staff default: Key on the transactional boundary entity. Accept and name the fan-out queries; serve them from CDC-fed indexes.

When to deviate: If the dominant workload is analytical range scans over time (logs, metrics), key on time buckets. If no single entity bounds transactions, stop hand-sharding.

🧭 Principal Move: "The shard key is an org-wide contract, not a table property. If orders, payments and fulfillment all shard by merchant_id, a merchant's data co-locates across services and cell-based isolation becomes possible. If each team picks its own key, every cross-service workflow is cross-shard forever."

3.2 Fault Line 2: Placement — Algorithmic vs Directory#

The tension: Algorithmic placement (hash → shard) has no dependency and never goes stale, but can't place individual entities. A directory can put any entity anywhere, but becomes a critical dependency and a source of stale routing.

ApproachWhat WorksWhat BreaksWho Pays
hash % NZero stateChanging N remaps nearly everythingWhoever reshards (full rewrite)
Consistent hashing ringAdding node moves ~1/NCan't pin a whale; ring changes still move TBsDB platform
Per-entity directoryArbitrary placementDirectory lookup per request; directory outage = total outageEvery request (latency) + directory owner
Hash → logical shard → physical map (+ small override table)Rebalancing by moving logical shards; whales pinnedMap is a one-way door on logical count; override table must stay smallDB platform (owns map)

Staff default: Two-level: logical = hash(key) % 4096, physical = map[logical], plus an override table limited to ~1,000 entries. The full map is ~4K entries (~100 KB) — every router caches all of it. Updates are versioned; routers reject requests stamped with a stale version from a shard that no longer owns the data.

The stale-map problem: During a move, a router with map v416 sends a write to the old cluster after cutover to v417. Defense: each shard stores the set of logical shards it owns and the map version; it rejects writes for logical shards it no longer owns with WRONG_SHARD(v417), and the router refreshes and retries. This is fencing at the data layer, and it's what makes directory placement safe.

Diagram: 3.2 Fault Line 2: Placement — Algorithmic vs Directory

3.3 Fault Line 3: Resharding — Planned Online Migration vs Emergency Split#

The tension: Building copy/catch-up/verify/cutover tooling costs 1–2 engineer-quarters before you need it. Not building it means your first resharding happens when a shard is at 95% disk, under incident pressure, with hand-written scripts.

The online migration workflow (per logical shard):

Diagram: 3.3 Fault Line 3: Resharding — Planned Online Migration vs Emergency Split
PhaseWhat HappensDuration (500 GB shard)Rollback
CopyConsistent snapshot of logical shard to target; record WAL/binlog position45–170 min at 50–200 MB/sDrop target
Catch-upReplay changes since snapshot positionUntil lag < 1 sDrop target
VerifyChunked checksums per PK range, source vs target20–60 minRepair chunks, re-verify
FreezeBlock writes for this logical shard at the router1–5 sUnfreeze
FlipPublish map v+1; shards fence on version< 1 sPublish v back
Reverse replicationNew → old stream keeps old current24–72 hFlip back — no data loss
CleanupDelete moved rows from sourceThrottled, hours—

Who pays: Without tooling, the on-call pays during the incident and customers pay in minutes-to-hours of write unavailability. With tooling, the DB platform team pays 1–2 quarters upfront plus ~10% ongoing.

Staff default: Build the workflow (or adopt Vitess / a managed engine that has it) before the first shard exceeds 60% capacity. Run it quarterly as a game day even when not needed, so the path is exercised.

When to deviate: If growth is predictable and slow (< 30%/year), over-provision logical shards and physical hosts up front and accept rare, planned maintenance-window moves.

🎯 Staff Move: "Resharding isn't an event, it's a capability. I want a runbook that moves one logical shard with a 5-second write pause, a verify step that blocks cutover on any checksum mismatch, and reverse replication so rollback is a config flip, not a restore."

3.4 Fault Line 4: Cross-Shard Correctness — Distributed Transactions vs Sagas#

The tension: Atomic commit across shards gives simple invariants but couples availability (all participants must be up) and holds locks across round-trips. Sagas keep shards independent but expose intermediate states.

ApproachWhat WorksWhat BreaksWho Pays
2PC / XAAtomicityCoordinator failure after prepare blocks participants; +5–20 ms; lock contentionAll shards in the txn (coupled availability)
Distributed SQL (Spanner-style)Serializable across shardsCommit-wait / consensus latency (≈10–100 ms cross-region); costLatency budget + infra spend
Saga + outboxShards stay independent; each step localIntermediate states visible; compensations must be correctProduct (UX for pending) + service owner (compensation logic)
Design it awayNo cross-shard txnSome entity relationships denormalizedData modeling effort up front

Staff default: Design away ≥95% of cross-shard writes via the key. For the rest: transactional outbox on the initiating shard + idempotent saga steps keyed by a business ID + compensation + a reconciliation job. Name the drifting invariant and its bound.

Diagram: 3.4 Fault Line 4: Cross-Shard Correctness — Distributed Transactions vs Sagas

When to deviate: Financial ledgers where regulators require atomic double-entry across accounts — co-locate the ledger by account group, or use a distributed SQL engine for that one table.

3.5 Fault Line 5: Routing Ownership — Library vs Proxy vs Engine#

OptionWhat WorksWhat BreaksWho Pays
Routing library in each serviceNo extra hopVersion skew across 3–4 languages; map refresh bugs per languageEvery product team
Sharding proxy (Vitess, ProxySQL, Citus coordinator)One implementation; online resharding toolingExtra hop (0.3–1 ms); proxy fleet to runDB platform team
Distributed SQL engineSharding invisible to appsHigher per-query latency; migration effort; vendor lock-inBudget + latency
Managed KV (DynamoDB)Zero sharding opsData model constrained to key-value/GSI patternsProduct teams (modeling limits)

Staff default: Proxy for polyglot orgs with > 3 services on the sharded DB; library only when one service owns the data. For new systems without heavy relational needs, a managed store with built-in partitioning.

🎯 Staff Move: "The routing layer is a platform product. If each team embeds its own router, the first map-format change requires coordinated deploys across every service — that's the SDK hell we already know from rate limiting."


4. Failure Modes & Operational Reality#

4.1 Hot Shard — The Whale Merchant#

Scenario: A single merchant runs a flash sale. Their logical shard lives on Cluster 07 alongside ~255 other logical shards.

t=0:       Merchant 88213 launches a flash sale; 12K orders/min on one merchant_id
t=+30s:    Cluster 07 primary CPU 45% → 92%; p99 write latency 8 ms → 340 ms
t=+1min:   255 unrelated merchants on Cluster 07 see checkout timeouts
t=+2min:   Connection pool saturation; app retries double the load
t=+3min:   Replica lag 30 s; merchant dashboards on Cluster 07 show stale orders
t=+8min:   On-call identifies merchant 88213 via per-key top-N
t=+10min:  Rate limit on merchant 88213 checkout (degraded mode) — neighbors recover
t=+50min:  Logical shard for 88213 moved to a dedicated cluster via override

Detection: shard.cpu_pct{cluster} > 80% for 2 min; shard.p99_write_ms > 50; shard.top_keys_qps showing one key > 20% of the shard. Blast radius: Every tenant co-located on the cluster — ~1/16 of merchants. Mitigation: Per-key admission control first (seconds), then move the logical shard (tens of minutes). Prevention: Directory overrides for the top 0.1% of merchants by GMV; per-key QPS budgets; pre-event capacity reviews for known campaigns. Owner: DB platform on-call (move); checkout team (per-merchant limits).

🎯 Staff Move: "Adding shards doesn't help a hot key — it spreads the neighbors, not the whale. The fix is isolation: the whale gets its own cluster, and its neighbors stop paying for it."

4.2 Stale Shard Map — Writes to the Wrong Place#

t=0:       Migration flips logical 1207 from Cluster 03 to Cluster 19 (map v417)
t=+0.2s:   38 of 40 routers pick up v417 via watch
t=+0.2s:   2 routers had a broken etcd watch after a network blip; still on v416
t=+0→45s:  Those 2 routers write ~1,100 orders to Cluster 03 for logical 1207
t=+45s:    Periodic 30 s poll... also failing (same connection issue)
t=+6h:     Customer reports missing order; reconciliation finds 1,100 orphaned rows

Root cause: Correctness depended on every router having the latest map. Prevention: shard-side fencing — Cluster 03 knows it no longer owns 1207 after the flip and rejects the write with WRONG_SHARD. With fencing, the same incident is 45 seconds of retries, not 1,100 orphaned rows. Detection: router.map_version_skew (max − min across fleet) > 0 for > 60 s; shard.wrong_shard_rejections_total. Owner: DB platform.

4.3 Resharding Gone Wrong — The Silent Diff#

Scenario: Catch-up replication skips rows written by a bulk job that bypassed the binlog (e.g., LOAD DATA with replication filters, or an unlogged table).

t=0:       Copy of logical shards 512–767 to new cluster begins
t=+2h:     Catch-up lag < 1 s; operator skips verify "to save time"
t=+2h05m:  Cutover
t=+3d:     Finance: 0.3% of refunds from a nightly batch are missing on the new cluster
t=+3d:     Reverse replication already torn down; old cluster cleaned up at t+48h

Lesson: Verification is not optional, and cleanup must wait for a business-cycle soak (at least one full nightly/weekly batch cycle). Detection: migration.verify_mismatch_chunks must be 0 to proceed — enforced by the orchestrator, not the operator. Owner: DB platform (tooling guardrail); batch job owners (must write through normal paths).

4.4 Fan-Out Tail Latency Collapse#

A "search orders by status" admin feature is implemented as scatter-gather over 64 clusters. Each cluster's p99 is 20 ms, but p99 of the fan-out is the p99.98 of individual shards.

Shards queriedP(at least one shard > own p99)Effective latency
11%20 ms p99
1615%~p99 of shard at p85 of requests
6447%Half of requests hit a slow shard
25692%Nearly every request is slow

Worse: at 200 admin QPS × 64 shards = 12,800 shard queries/s of full-index scans — a load amplifier that hurts OLTP on every shard. Detection: router.fanout_ratio (shard-queries per client query) > 1.2 fleet-wide; per-endpoint fan-out breakdown. Mitigation: Serve cross-shard queries from a CDC-fed search index or warehouse; cap fan-out QPS via a separate, rate-limited pool. Owner: The feature team that introduced the fan-out; DB platform sets the budget.

4.5 Cross-Shard Saga Stuck in the Middle#

The outbox relay crashes mid-deploy; 40K orders are PENDING with credit already debited on another shard. Detection: saga.pending_age_seconds p99 > 60 s; outbox.unpublished_rows growing. Mitigation: Relay is idempotent; restart drains the backlog. Orders older than 15 min go to a compensation job that re-checks the credit shard. Owner: Checkout team (saga logic); platform (relay infra).

4.6 Operational Reality Matrix#

FailureDetection SignalBlast RadiusMitigationOwner
Hot shard (whale)shard.cpu_pct > 80%, shard.top_keys_qpsAll tenants on that clusterPer-key limits → directory override moveDB platform + product
Shard primary failureshard.primary_up = 01/N of entities, writes onlyAutomated failover to replica (~10–30 s)DB platform
Stale shard maprouter.map_version_skew > 0 for 60 sKeys in moved logical shardsShard-side fencing; forced refreshDB platform
Resharding diffmigration.verify_mismatch_chunks > 0Moved logical shardsBlock cutover; repair chunksDB platform
Fan-out overloadrouter.fanout_ratio > 1.2All shardsMove query to derived store; cap fan-out poolFeature team
GSI laggsi.lag_seconds p99 > 5Buyer views, searchScale CDC consumers; degrade UI with "may take a moment"Data platform
Saga stucksaga.pending_age_seconds > 60In-flight cross-shard opsIdempotent replay; compensation jobOwning product team
Directory (etcd) downdirectory.available = 0Moves blocked; routing continues from cacheFreeze migrations; routers serve cached mapDB platform
Shard disk fullshard.disk_pct > 85%That cluster's writesEmergency move of largest logical shards; archivalDB platform

5. Evaluation Rubric#

5.1 Level-Based Signals#

DimensionSenior (L5)Staff (L6)Principal (L7)
Need to shardAssumes sharding is requiredChecks alternatives (archival, partitioning, replicas) and names the ceiling with numbersPrices build (sharded MySQL/Postgres) vs buy (distributed SQL / managed KV) over 3 years
Shard keyPicks a reasonable IDDerives from transactional boundary; lists fan-out queries and how each is servedMakes the key a cross-service convention; aligns with cell architecture
PlacementConsistent hashingLogical shards + directory + overrides + data-layer fencingSets logical shard count as a documented one-way door with a 5-year growth model
Resharding"Rebalance"Copy/catch-up/verify/freeze/flip/reverse-replicate with numbersResharding as a platform SLO; quarterly game days; tooling reused org-wide
Cross-shard2PCDesigns away; outbox + saga; names drifting invariant and boundOrg policy on which invariants may be eventual; audit + reconciliation ownership
OperationsReplicas for HAPer-shard metrics, fencing, fan-out budget, ownership matrixBlast-radius cells, error budgets per cluster, cost per shard in the capacity plan

5.2 Strong Hire Signals#

SignalWhat It Sounds Like
Query-first key selection"Checkout is atomic per merchant, so merchant_id. Buyer history fans out, so it becomes an async GSI with 1–2 s lag."
Resharding as a workflow"Snapshot, binlog catch-up, chunked checksum verify, 3-second freeze, map flip, 48-hour reverse replication."
Fencing awareness"Routers can have stale maps, so each shard rejects writes for logical shards it doesn't own."
Hot shard vs hot key"More shards won't help a single hot merchant. Isolate it; split its key only if one merchant exceeds a cluster."
Named drifting invariant"Credit balance vs orders can drift for up to 5 seconds; hourly reconciliation; finance signed off."

5.3 Lean No-Hire Signals#

SignalWhy It Misses the Bar
Hash mod NRemaps nearly all keys on any change — no resharding story
2PC everywhereCouples availability of all shards; no awareness of coordinator failure blocking
"Scatter-gather is fine"Ignores tail amplification and load amplification
No mention of migrating the live systemGreenfield-only thinking; the real project is moving from the monolith
Shard key = auto-increment IDEither range-hot or destroys entity locality

5.4 Common False Positives#

  • Deep knowledge of consistent-hashing math ≠ sharding judgment. Virtual-node counts don't answer "which queries fan out?"
  • Naming Vitess/Citus ≠ understanding. Can they describe what MoveTables actually does and where it can fail?
  • Proposing Spanner immediately can be correct — but only if they price the latency and cost and explain why hand-sharding loses.
  • Elaborate GSI designs — synchronous global indexes are a distributed transaction in disguise.

6. Interview Flow & Pivots#

6.1 Typical 45-Minute Shape#

PhaseTimeGoal
Framing0–3 minTop 3 queries, transactional boundary, write QPS and size numbers
Entities + API3–5 minKey-carrying API; IDs embed logical shard
High-level design5–10 minRouter, logical→physical map, clusters, CDC side paths
Transition10–11 minOffer shard-key defense / resharding / cross-shard
Deep dives11–40 minResharding workflow → hot shards → cross-shard sagas
Wrap-up40–45 minEvolution, what not to build, ownership

6.2 How Interviewers Pivot — And What They're Testing#

Interviewer SaysWhat They're TestingWhere to Go
"Now the business wants to query by buyer"Do you know you made it a fan-out?Async GSI via CDC, lag budget
"We need to double capacity next quarter"Resharding maturityLogical shard moves, verify, fencing
"One tenant is 30% of traffic"Hot key vs hot shardDirectory override → intra-tenant key
"Transfer money between two users"Cross-shard correctnessOutbox + saga, drifting invariant
"Why not just use Spanner/DynamoDB?"Build vs buyTCO, latency, data model constraints
"We're moving to 3 regions"Geo placementHome region per entity, residency

6.3 What to Deliberately Skip#

TopicWhy L5 Goes HereWhat L6 Says Instead
Ring math, virtual nodesFeels rigorous"Logical shards give the same effect with explicit control."
Replication within a shardFamiliar"Standard primary + 2 replicas per cluster."
Query planner designSounds hard"Cross-shard analytics live in the warehouse."
ID generator internalsFun"64-bit, time-ordered, logical shard embedded. Instagram-style."

6.4 Follow-Up Questions to Expect#

  1. "How do you pick the number of logical shards, and what if you pick wrong?"
  2. "Walk me through adding capacity with zero write loss."
  3. "What happens to in-flight transactions during cutover?"
  4. "How do you enforce a uniqueness constraint (e.g., unique email) across shards?"
  5. "A migration corrupts data on the target. How do you find out, and how do you roll back?"
  6. "How do you run a schema migration across 64 clusters?"
  7. "How would you shard if the product later requires multi-merchant carts?"

7. Active Drills#

Drill 1: The Opening#

Prompt: "Our Postgres primary is at 85% CPU. Design sharding for it."

Staff Answer

"First, I want to confirm sharding is the right lever. 85% CPU — is it writes or reads? If reads, replicas and caching are cheaper. If it's vacuum or index bloat on one large table, partitioning and archival may buy a year. Assuming it's genuine write throughput — say 25K writes/s growing 2× yearly — I'll shard.

Next, the transactional boundary. For a marketplace, checkout touches one merchant's inventory and creates an order, so merchant_id is the key. I'll hash into 4,096 logical shards mapped onto 8 clusters initially, with a directory override for the top merchants. Buyer history and admin search become CDC-fed derived stores with 1–2 seconds of lag. The real project is migrating the live monolith, so I'll design the online migration first — the same tooling becomes our resharding capability."

Why this is L6:

  • Challenges the premise with specific cheaper alternatives before committing
  • Chooses the key from the transactional boundary and names what fans out
  • Recognizes the migration from the monolith as the real risk

What L7 adds:

  • Prices hand-sharding (≈2–3 engineers ongoing + tooling) against a managed distributed SQL engine and states the break-even
  • Sets merchant_id as the org-wide partition convention so payments and fulfillment co-locate
  • Frames the migration tooling as a reusable platform capability for other teams

Drill 2: Choosing the Logical Shard Count#

Prompt: "Why 4,096 logical shards? Why not 16 or a million?"

Staff Answer

"Logical shards are the unit of movement, so the count bounds both granularity and overhead. Too few (16): each move is 1/16 of the data, and you can never grow beyond 16 hosts without rehashing. Too many (1M): the map is large, per-shard metadata (schemas, connection pools, monitoring series) explodes, and moves become death-by-a-million-cuts.

With 12 TB today at 2× growth, in 5 years we're at ~380 TB. At a 500 GB target per host that's ~760 hosts — 4,096 logical shards gives ~5 per host at that point, still enough to balance. Today, on 8 hosts, each holds 512 logical shards of ~6 GB — each moves in about a minute. If I wanted cleaner divisibility across host counts, I'd consider a highly composite number like 480 or 5,040; powers of two are fine if host counts are also powers of two."

Why this is L6:

  • Grounds the number in growth math, not taste
  • Explains both failure directions (too few, too many)
  • Connects the count to move duration

What L7 adds:

  • Declares the count a one-way door in the design doc, with the re-hash cost estimated (a full re-key migration ≈ 1–2 quarters)
  • Sets a review trigger: "revisit when average logical shard exceeds 50 GB"

Drill 3: Make Resharding Concrete#

Prompt: "You're moving from 8 clusters to 16. Walk me through it."

Staff Answer

"Placement change only: move 2,048 of 4,096 logical shards — half of each source cluster — to the 8 new ones. No rehashing. I'd batch in waves of ~64 logical shards.

Per wave: take a consistent snapshot of each logical shard (per-schema dump or a filtered snapshot), record the binlog/WAL position, copy at ~100 MB/s, then stream changes from that position until lag < 1 s. Verify with chunked checksums over PK ranges; any mismatch blocks cutover. Cutover: router pauses writes for those logical shards (~2 s), waits for zero lag, publishes map v+1; source clusters fence on version so stale routers get WRONG_SHARD and refresh. Reverse replication from new to old runs for 48 hours — rollback is republishing the old map. Then throttled delete on the source.

Guardrails: max 2 waves concurrently, pause automatically if shard.p99_write_ms on any source exceeds 2× baseline, and no waves during peak hours. Total: about 2 weeks elapsed for 12 TB, mostly soak time."

Why this is L6:

  • Every phase has a number and a rollback
  • Fencing handles stale routers without depending on perfect propagation
  • Load-aware throttling protects production during the move

What L7 adds:

  • Makes the orchestrator self-service with automated guardrails, so moves don't need a Staff engineer on call
  • Tracks resharding.moves_per_quarter and resharding.rollback_rate as platform health KPIs

Drill 4: The Directory Goes Down#

Prompt: "Your etcd cluster holding the shard map is down for 20 minutes. What breaks?"

Staff Answer

"If designed correctly: nothing on the data path. Routers hold the full map (~100 KB) in memory and on local disk; they keep routing with the last known version. What breaks: moves can't cut over, new routers can't bootstrap unless they load from a snapshot in object storage, and override changes are frozen.

The orchestrator must detect directory unavailability and pause any wave in the Freeze state — unfreeze and abort, don't hold writes paused for 20 minutes. Alert directory.available == 0 pages DB platform; router.map_age_seconds tells us how stale the fleet is. The failure I'd worry about is a router that restarts during the outage — so routers boot from a signed snapshot of the map published to S3 on every version change."

Why this is L6:

  • Keeps the control plane off the data path
  • Handles the in-flight migration explicitly
  • Covers cold-start routers

What L7 adds:

  • Makes "control plane can be down for 1 hour with no data-plane impact" a written requirement for every platform service
  • Exercises it in a quarterly game day

Drill 5: Hot Key — One Merchant at 30% of Writes#

Prompt: "One merchant is now 30% of all order writes. What do you do?"

Staff Answer

"30% of 40K writes/s is 12K writes/s on one key — more than a comfortable single cluster. Step 1, within minutes: per-merchant admission control so neighbors aren't hurt. Step 2, within an hour: directory override to a dedicated cluster. Step 3, if one cluster isn't enough: introduce a second-level key inside this merchant — (merchant_id, store_id) or (merchant_id, hash(order_id) % 8) — which means this merchant's inventory must be partitioned by warehouse/store so checkout stays single-shard. That's a data-model change the merchant's product surface must accept.

If the hot spot is a single row — e.g., one SKU's inventory counter — no sharding helps. I'd split the counter into N sub-counters or use reservation batches."

Why this is L6:

  • Distinguishes hot shard, hot tenant and hot row
  • Stages the response by time horizon
  • Names the data-model consequence of splitting a tenant

What L7 adds:

  • Turns whale handling into contract tiers: enterprise merchants above X GMV get a dedicated cell, priced into their plan
  • Defines the whale-detection threshold and ownership in the platform standard

Drill 6: Global Uniqueness#

Prompt: "Emails must be unique across all users, but users are sharded by user_id."

Staff Answer

"A unique constraint on a non-key column is a cross-shard invariant. Options: (1) a separate uniqueness table sharded by hash(email) — signup writes a claim there first, then the user row, with the claim as the source of truth; a background job cleans claims without a user row after 10 minutes. (2) Fan-out check — racy and slow, no. (3) Distributed transaction — overkill for signup.

I'd pick (1): the email-claim table is itself a sharded table keyed by email, so the uniqueness check is a single-shard insert with a unique index. The only cross-shard step is claim → user row, handled as a two-step idempotent flow. Signups are ~100/s — nowhere near a scaling concern."

Why this is L6:

  • Recognizes uniqueness as a derived index keyed on the unique column
  • Orders the writes so failures leave cleanable garbage, never duplicates

What L7 adds:

  • Publishes this "claim table" pattern as the standard for all global-uniqueness needs (usernames, handles, coupon codes)

Drill 7: Build vs Buy#

Prompt: "Why not just use a distributed SQL database instead of sharding Postgres ourselves?"

Staff Answer

"It's a legitimate option and sometimes the right one. Distributed SQL gives automatic range splitting, online rebalancing and cross-shard transactions. The costs: single-row write latency is typically higher (consensus round-trips, ~2–10 ms in-region vs ~1 ms for a local Postgres commit), operational expertise is different, and migrating 12 TB is itself a project.

I'd buy if: our transactions genuinely span entities, our team has < 2 engineers who could own sharding tooling, or we need multi-region writes. I'd build (or adopt Vitess/Citus) if 95% of transactions are single-entity, we have deep Postgres/MySQL expertise, and latency matters at p99."

Why this is L6:

  • Treats buying as a real option with specific triggers
  • Quantifies the latency difference

What L7 adds:

  • Frames the decision over 3 years including headcount (hand-sharding ≈ 3–5 FTE platform team at scale) and exit cost
  • Negotiates the vendor relationship: data export guarantees, pricing at 10× scale

Drill 8: Schema Change Across 64 Clusters#

Prompt: "Add a column to orders. You have 64 clusters."

Staff Answer

"Schema changes become a fleet rollout. Use online DDL tooling (gh-ost/pt-osc for MySQL, or Postgres' fast ADD COLUMN with defaults) orchestrated as a canary: 1 cluster → 4 → 25% → 100%, with automatic pause on replication lag > 10 s or p99 regression > 20%. Application code must tolerate both schemas (expand/contract): deploy readers that ignore the column, add the column everywhere, deploy writers, backfill, then enforce NOT NULL. A partially migrated fleet is the normal state for hours — design for it."

Why this is L6:

  • Treats DDL as a staged rollout with automatic guardrails
  • Uses expand/contract so mixed-schema fleets are safe

What L7 adds:

  • Centralizes schema-change orchestration so 30 teams don't each script it; tracks schema.fleet_convergence_hours

Drill 9: Cost#

Prompt: "Finance says the database bill doubled after sharding. Defend it or fix it."

Staff Answer

"Sharding does add overhead: 16 clusters × 3 replicas is 48 nodes versus 3 before, plus routers and CDC. First, check utilization: if the average cluster is at 25% CPU, we over-provisioned — consolidate logical shards onto fewer hosts (that's what logical shards are for). Second, right-size replicas: do all clusters need 2 read replicas, or only the hot ones? Third, move cold data (> 18 months) to object storage — often 50–70% of bytes. Target: average cluster at 50–60% of peak capacity with a 2× surge buffer. I'd expect to cut the increase by 30–40% without touching reliability."

Why this is L6:

  • Uses the design's flexibility (logical shards) to fix cost
  • Gives a utilization target, not "optimize"

What L7 adds:

  • Introduces per-team showback (cost per logical shard, attributed by tenant) so product teams see their data's cost

Drill 10: Multi-Region#

Prompt: "We're expanding to the EU and data must stay in-region."

Staff Answer

"Add region as the top level of placement: each merchant has a home region recorded in the directory. EU merchants' logical shards live on EU clusters; routers in each region route locally, and cross-region requests for an EU merchant are forwarded to the EU router, never replicated out. Global tables (e.g., the email-claim table) must be split or made region-aware; the warehouse needs region-scoped pipelines. Moving an existing merchant to the EU is the same online-migration workflow, cross-region, with a longer catch-up (latency ~80–150 ms)."

Why this is L6:

  • Reuses the directory as the residency mechanism
  • Identifies global tables as the hidden problem

What L7 adds:

  • Makes region an attribute in the entity model org-wide; legal signs off on the data classification per table

8. Deep Dive Scenarios#

Deep Dive 1: Peak-Traffic Hot Shard#

Context: Singles' Day. Cluster 07 is at 98% CPU, checkout p99 is 2 s for ~6% of merchants, and the on-call escalates to you. One merchant is running a live-streamed flash sale.

Questions to Surface First:

  • Is it one merchant (hot key) or many busy merchants on one cluster (hot shard)?
  • Are neighbors' errors due to CPU, locks, or connection-pool exhaustion?
  • Is the merchant's traffic legitimate, and what's the revenue at stake for them vs neighbors?
  • Is a move possible right now without making the source cluster worse?

Typical L5 Approach: Scales the cluster vertically or adds read replicas, then proposes adding more shards. Neither helps within the hour: vertical scaling requires failover, replicas don't absorb writes, and more shards don't split one merchant.

Staff Approach: Identify the key via top-N. Apply per-merchant admission control to protect neighbors immediately (queue the whale's checkouts, don't fail them). Then decide: moving under load adds copy I/O to a saturated primary — throttle the copy to 30 MB/s or wait until the peak passes. After the event, pre-place known whales before campaigns.

Principal Approach: The incident is a symptom of missing tenant tiers. Large merchants get dedicated cells by contract; campaign calendars feed capacity planning; the platform publishes "per-tenant write budget" as part of the merchant API contract. The question for the VP is whether we charge for isolation or absorb it.

Staff Approach — Full Reasoning
PhaseWhat to Do
Immediate (0–5 min)shard.top_keys_qps → confirm merchant 88213 is 70% of cluster writes. Enable per-merchant checkout concurrency cap (e.g., 400 in-flight).
TriageNeighbors recover? Check shard.p99_write_ms for other merchants. Whale gets queued checkout with visible "processing" state.
Quick fixAfter peak: throttled move of the whale's logical shard to a pre-provisioned spare cluster via directory override.
GuardrailsPause move if source p99 > 2× baseline; keep reverse replication for 48 h.
Post-mortemWhy wasn't the campaign in capacity planning? Add whale detection (> 5% of cluster writes for 1 h) that auto-files a placement ticket.

Metrics to Watch: shard.cpu_pct, shard.top_keys_qps, checkout.queue_depth{merchant}, shard.lock_wait_ms, migration.copy_mb_s

Organizational Follow-up: Marketing/merchant-success must file large campaigns 7 days ahead; DB platform pre-places. Merchant tier contract includes dedicated-cluster eligibility.

Ownership Question: "Who decides to throttle a top-10 merchant during their biggest sale?" Staff answer: The checkout on-call, under a pre-approved runbook with merchant-success notified — the policy (protect the fleet over one tenant, queue rather than reject) was signed off by the business before the event, not debated during it.

Key Takeaway: "More shards spread neighbors, not whales. Isolation is the fix; admission control buys the time to do it."

What clears the Staff bar:

  • Separates hot key from hot shard within the first minute
  • Protects neighbors before fixing the whale
  • Converts the incident into a placement policy

Deep Dive 2: Silent Data Divergence After Resharding#

Context: Three weeks after a resharding wave, a support ticket reveals an order that exists in the warehouse but not in the OLTP database. Reconciliation finds ~2,300 missing rows across 14 logical shards.

Questions to Surface First:

  • Were the rows written before, during or after cutover?
  • Did the verify step run and pass for those logical shards?
  • Is there any write path that bypasses the router (batch jobs, direct DB access, triggers)?
  • Is the old cluster's data still available?

Typical L5 Approach: Restores from backup and re-inserts the rows. Fixes the symptom; the next wave loses rows the same way.

Staff Approach: Treat it as a correctness incident. Timestamp analysis shows all rows came from a nightly refund job connecting directly to the old cluster by hostname — bypassing the router, so fencing never saw them. Fix: backfill from the old cluster (still retained per policy) and the warehouse; revoke direct DB credentials; enforce router-only access with network policy; add per-logical-shard row-count reconciliation between OLTP and warehouse daily.

Principal Approach: Direct database access by hostname is an org-level anti-pattern that will break every future platform migration. Write the standard: all data access goes through the routing layer; exceptions need a registered, fenced client. Audit all 40 services' connection strings; publish a deadline. The fix is governance, not one job.

Staff Approach — Full Reasoning
PhaseWhat to Do
Immediate (0–5 min)Freeze further resharding waves. Confirm old clusters' data retention hasn't expired.
TriageDiff by logical shard and write timestamp; correlate with client identity in DB audit logs.
Quick fixReplay missing rows from old cluster into correct shard, idempotently by PK. Notify affected merchants' finance teams if refunds were delayed.
GuardrailsOld cluster rejects all writes for moved logical shards regardless of client (fencing at the DB, not the router).
Post-mortemWhy could a job bypass the router? Why didn't verify catch it? (It ran before the nightly job wrote.) Extend soak to cover a full weekly batch cycle.

Metrics to Watch: reconcile.row_count_diff{logical_shard}, shard.writes_by_client, shard.wrong_shard_rejections_total

Organizational Follow-up: Inventory all direct-DB clients; move batch jobs onto the router; add "fenced after cutover" as a migration exit criterion.

Ownership Question: "Who owns the missing refunds?" Staff answer: DB platform owns the data repair and the guardrail; the refund job's team owns the customer impact and fixing their access path; finance signs off on the reconciliation result.

Key Takeaway: "Fencing only works if every writer goes through the fence. A migration's blast radius is every client you forgot about."

What clears the Staff bar:

  • Looks for bypass paths, not just tooling bugs
  • Moves fencing to the storage layer
  • Extends soak to cover business cycles

Deep Dive 3: Large-Customer Onboarding#

Context: Sales closes an enterprise customer migrating 40M orders and expecting 3K writes/s at launch — larger than any current merchant. Launch is in 6 weeks.

Questions to Surface First:

  • Does this customer fit on one cluster with 2× headroom?
  • Do they need data residency or contractual isolation?
  • How is the 40M-order import staged — bulk load or through the API?
  • What's their 12-month growth forecast?

Typical L5 Approach: Imports through the normal API over a weekend. 40M orders at API rate limits takes days and hammers the shared cluster.

Staff Approach: Pre-provision a dedicated cluster; directory override for the merchant before import; bulk-load directly into the dedicated cluster via the migration tooling's copy path (not through OLTP APIs); build GSIs and warehouse via CDC backfill in a separate, throttled pipeline. Load-test at 2× expected peak (6K writes/s) two weeks before launch.

Principal Approach: Enterprise onboarding is a repeatable product: an "onboarding cell" template with pricing, a capacity SLA, and a lead time. Sales gets a checklist that triggers platform involvement at contract stage, not after signing.

Staff Approach — Full Reasoning
PhaseWhat to Do
Week 1Capacity model: 3K writes/s + 40M rows (~120 GB) fits one cluster at ~40% CPU. Provision dedicated cluster + override.
Week 2–3Dry-run bulk import into staging; measure copy rate and GSI backfill lag.
Week 4Production import; GSI/warehouse backfill throttled to keep gsi.lag_seconds for other merchants < 5 s.
Week 5Load test at 6K writes/s; failover drill on the dedicated cluster.
LaunchDedicated on-call coverage for 72 h; rollback = pause the merchant's API keys, not the fleet.

Metrics to Watch: shard.cpu_pct{cluster=dedicated}, gsi.lag_seconds, import.rows_per_s, checkout.p99_ms{merchant}

Organizational Follow-up: Sales-to-platform intake for any deal above 1K writes/s or 50 GB.

Ownership Question: "Who pays for the dedicated cluster?" Staff answer: The customer, via the enterprise tier price. If sales discounts it away, the business has explicitly chosen to subsidize isolation — that's their call, but it must be visible.

Key Takeaway: "Onboard whales by placement, not by hope. The directory exists so the biggest customer never shares a cluster by accident."

What clears the Staff bar:

  • Bulk path separate from OLTP path
  • Isolates derived-store backfill from other tenants
  • Connects capacity to commercial terms

Deep Dive 4: Post-Mortem — The 2PC Outage#

Context: A previous team implemented cross-shard transfers with XA/2PC. The coordinator host lost its disk after participants prepared. 11 minutes of locked rows on 6 clusters cascaded into a checkout outage. You're asked to lead the post-mortem.

Questions to Surface First:

  • How many transactions were in-doubt, and which rows did they lock?
  • Did the coordinator's decision log survive?
  • What fraction of all writes are cross-shard?
  • Why did locks on transfer rows block checkout?

Typical L5 Approach: Make the coordinator highly available (replicate its log). Reduces probability but keeps the coupling: any participant slowness still holds locks across shards.

Staff Approach: Measure: cross-shard transfers were 0.8% of writes but caused 100% of this outage. Replace 2PC with outbox + saga: debit on source shard with an outbox entry, idempotent credit on target, compensation on failure, hourly ledger reconciliation. Invariant "sum of balances is constant" is violated by in-flight transfers for ≤ 5 s — finance signs off because the ledger records both legs with a transfer ID.

Principal Approach: Publish a policy: synchronous cross-shard transactions require an exception review by the data platform council. Provide a saga/outbox library as the paved road so teams don't reinvent it. Track crossshard.sync_txn_count fleet-wide as a debt metric.

Staff Approach — Full Reasoning
PhaseWhat to Do
ImmediateManually resolve in-doubt transactions from participants' prepared state (XA RECOVER); commit or roll back per business rules.
TriageMap lock dependencies: transfer rows shared an index with checkout's balance reads.
Quick fixDisable transfers (feature flag) until saga rollout; ~0.8% of volume, acceptable for 2 weeks.
GuardrailsLock timeouts ≤ 2 s on all clusters; alert on prepared transactions older than 5 s.
Post-mortemRoot cause: design coupled availability of 6 clusters to 1 coordinator. Contributing: no lock-timeout policy.

Metrics to Watch: db.prepared_txn_age_seconds, db.lock_wait_ms, saga.pending_age_seconds, ledger.reconcile_diff

Organizational Follow-up: Saga library owned by platform; transfer team owns compensation logic; finance owns reconciliation sign-off.

Ownership Question: "Who approves re-enabling transfers?" Staff answer: The transfer team's lead plus finance, after a shadow-run of the saga matches the old ledger for 7 days with zero reconciliation diffs.

Key Takeaway: "2PC turns the availability of N shards into the product of their availabilities — and the coordinator is a single point of blocking."

What clears the Staff bar:

  • Quantifies share of writes vs share of risk
  • Names the drifting invariant and who signs off
  • Converts one fix into a platform default

Deep Dive 5: Multi-Region Expansion#

Context: Leadership wants active-active across US-East, US-West and EU within 12 months, with EU data residency. Today: single-region, 32 clusters.

Questions to Surface First:

  • Is the goal latency, availability, residency — or all three?
  • Can each entity have a home region (writes go home), or must any region accept writes for any entity?
  • What's the tolerance for cross-region write latency (~70–150 ms RTT)?
  • Which tables are global (uniqueness claims, config)?

Typical L5 Approach: Multi-master replication for all clusters in all regions with last-writer-wins. Conflicts silently corrupt balances and inventory.

Staff Approach: Home-region placement: the directory records each merchant's home region; writes route home (~100 ms extra for remote users on writes, local reads from async replicas with a "your recent writes" session pin). EU merchants' data never replicates out. Global tables are split by region or become region-prefixed. Regional failover moves merchants' home to another region via the same migration workflow — with RPO > 0 for async replication, stated explicitly.

Principal Approach: Active-active is a 3-year program, not a feature. Stage it: Year 0 residency (EU cells), Year 1 regional read replicas + home-region writes, Year 2 automated regional failover with documented RPO/RTO. Budget: ~2.2× infra cost and one dedicated team. Say no to multi-master for money tables.

Staff Approach — Full Reasoning
PhaseWhat to Do
Quarter 1Directory gains home_region; routers become region-aware; EU clusters provisioned.
Quarter 2Migrate EU merchants via online workflow; legal verifies residency.
Quarter 3Cross-region async replicas for read latency; session pinning for read-your-writes.
Quarter 4Regional failover drill: move 5% of merchants' home region, measure RTO and data loss window.
Ongoingreplication.cross_region_lag_s p99 < 2 s is the published RPO.

Metrics to Watch: router.cross_region_forward_ratio, replication.cross_region_lag_s, residency.violations_total (must be 0)

Organizational Follow-up: Data classification per table (global/regional/residency-bound) owned by each product team, audited by legal.

Ownership Question: "Who decides to fail a region over and accept up to 2 s of data loss?" Staff answer: The incident commander, per a pre-agreed RPO that the business signed when funding the program — never an improvised decision.

Key Takeaway: "Give every entity a home. Multi-region sharding is just placement with a region column — until someone asks for multi-master."

What clears the Staff bar:

  • Separates latency, availability and residency goals
  • Reuses the directory and migration workflow
  • States RPO explicitly

9. Level Expectations Summary#

After studying this case study, you should be able to:

  • Decide whether to shard at all, using write QPS, size and growth numbers
  • Derive a shard key from the transactional boundary and list which queries fan out
  • Explain logical-shard placement with a directory, overrides and data-layer fencing
  • Walk through an online resharding workflow with durations, a write-freeze budget and rollback
  • Distinguish hot shard, hot tenant and hot row, with a mitigation for each
  • Replace distributed transactions with outbox + saga and name the drifting invariant
  • Price build vs buy and describe a 3-year evolution path

The Bar for This Question#

Mid-level (L4): Knows sharding splits data across nodes; can describe hash vs range partitioning and consistent hashing. Doesn't reason about resharding or cross-shard queries unprompted.

Senior (L5): Picks a sensible key, uses consistent hashing or logical shards, knows scatter-gather is slow, proposes 2PC or sagas when asked. Handles resharding at the "add nodes" level without a data-movement plan.

Staff+ (L6): Starts from queries and transactional boundaries, commits to a key and names what it costs. Designs the online migration with verify, fencing and rollback before being asked. Treats hot keys, cross-shard invariants and fan-out as owned operational risks with metrics. Knows when not to shard and when to buy. The interviewer should learn something from the answer.


10. Staff Insiders: Controversial Opinions#

10.1 "Most Teams Shard Two Years Too Early — and the Rest Shard Six Months Too Late"#

Premature shardingLate sharding
3 TB dataset, 5K writes/sPrimary at 90% CPU, largest instance
Adds router, directory, CDC, sagasMigration under incident pressure
Every feature pays the cross-shard taxHand-written scripts, no verify

The Staff position: Shard when your 18-month projection crosses ~60% of a single primary's ceiling — not earlier (you pay complexity for nothing), not later (you migrate in a fire).

Why this matters in interviews: Saying "we don't need to shard yet, here's the trigger" is a stronger signal than a perfect sharding design for a system that doesn't need one.

10.2 "Consistent Hashing Is the Wrong Answer to the Question You Were Asked"#

Consistent hashing minimizes how many keys move. Interviewers are asking how you move them safely. Logical shards achieve the same minimal movement with explicit, observable, controllable placement — and support pinning whales, which a ring can't.

The Staff position: Use consistent hashing for caches and Dynamo-style stores; use logical shards + a directory for relational OLTP.

Why this matters in interviews: Candidates who stop at "consistent hashing" haven't described the migration, which is the actual risk.

10.3 "Synchronous Global Secondary Indexes Are 2PC Wearing a Disguise"#

Updating a GSI on another shard in the same request means either a distributed transaction or an index that can be wrong after partial failure. DynamoDB's GSIs are asynchronous for exactly this reason.

The Staff position: GSIs are CDC-maintained, eventually consistent, with a published lag SLO (e.g., p99 < 2 s). Reads that need strong consistency go through the primary key.

10.4 "The Shard Key Is an Org Chart Decision"#

Once three services share a partition key, cells become possible and cross-service joins are local. Once they don't, you've built permanent cross-shard traffic into your architecture.

The Staff position: Choose the key with the adjacent teams in the room. It's cheaper to align now than to reconcile forever.

10.5 "If You Need Cross-Shard Transactions Everywhere, You Picked the Wrong Tool"#

Hand-sharding pays off when 95%+ of transactions are single-entity. If your domain is a graph of entities updated together, buy a distributed SQL engine and spend your engineers on the product.


11. The Principal Lens (L7)#

Why L7 Sees This Problem Differently#

A Staff engineer shards a database. A Principal engineer notices the company is about to shard its fifth database with its fifth bespoke router, its fifth migration script and its fifth incompatible shard key — and that the org will spend ~15 engineer-years maintaining them. The L7 question is not "what key for orders?" but "what is our partitioning platform, which entity is the company's partition axis, and which databases should never be hand-sharded at all?"

The Org-Level Fault Line#

One partitioning platform vs per-team sharding. Per-team sharding lets teams move fast and pick optimal keys locally; it produces divergent tooling, un-exercised migrations and cross-service data that never co-locates. A central platform (Vitess/Citus/managed distributed SQL + shared directory + migration orchestrator) standardizes operations but constrains data models and makes the platform team a bottleneck.

The L7 default: Centralize the mechanics (router, directory, migration, fencing, observability); decentralize the keys and schemas within a published convention (tenant/merchant as the company-wide partition axis). Offer a managed-KV path for teams whose model fits it.

Cost Model#

Assumptions: cloud list prices, 3-node clusters (primary + 2 replicas), ~$2.5K/month per node for 16 vCPU / 128 GB / 2 TB NVMe equivalent, fully loaded engineer $25K/month.

ScaleData / WritesInfra ($/month)Platform HeadcountOn-call LoadNotes
Small5 TB, 15K w/s, 8 clusters~$60K + routers/CDC $8K1–2 FTE (part-time)~2 pages/monthConsider not sharding; or managed distributed SQL at ~$80–120K
Medium50 TB, 100K w/s, 64 clusters~$480K + $40K4–6 FTE~8 pages/monthMigration tooling pays for itself; logical-shard balancing automated
Large500 TB, 1M w/s, 512 clusters~$3.8M + $250K10–15 FTEDedicated rotationCells + region placement; cost per logical shard in showback

Opportunity cost: at medium scale, a home-grown resharding tool is 2 engineer-quarters ($150K) — cheaper than one emergency resharding incident that burns a 6-person team for 3 weeks and causes a multi-hour checkout outage.

The 3-Year Evolution Path#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionDoorReversal Cost
Shard key (partition axis)One-wayFull re-key migration: 1–3 quarters, every dependent service
Logical shard countOne-way (mostly)Re-hash migration; mitigated by choosing 4,096+
IDs embedding logical shardOne-wayIDs are in URLs, logs, partner systems forever
Router: library vs proxyTwo-way1 quarter to swap with dual-run
Physical host count / instance typeTwo-wayMove logical shards; days
GSI technology (search vs KV)Two-wayRebuild from CDC; weeks
Saga vs 2PC for a flowTwo-way (painful)Rewrite one flow; weeks
Build (Vitess) vs buy (distributed SQL)One-way-ish12–24 month migration either direction

The Standard I'd Write#

RFC-DATA-007: Horizontal Partitioning of OLTP Stores

Scope: All OLTP datastores projected to exceed 2 TB or 10K writes/s within 18 months.

MUST:

  • Partition on the company partition axis (tenant_id/merchant_id) unless an exception is approved.
  • Use the platform router and directory; no direct host connections for application writes.
  • Hash into ≥ 1,024 logical shards; target ≤ 1 TB per physical cluster.
  • Enforce ownership fencing at the storage layer for every logical shard move.
  • Maintain secondary access paths via CDC with a published lag SLO.

SHOULD:

  • Keep ≥ 95% of write transactions single-shard; use the platform saga/outbox library otherwise.
  • Run one resharding game day per quarter.

Exceptions: Synchronous cross-shard transactions and non-axis keys require Data Platform Council review (48-hour SLA) with a named owner and a written invariant-drift analysis.

Success metrics: resharding.p50_move_hours < 2; crossshard.sync_txn_ratio < 1%; zero data-loss incidents from moves; ≥ 80% of eligible stores on the platform within 18 months.

What I'd Tell the VP#

"Our order database will hit its hardware ceiling in about 14 months at current growth. We can split it across many smaller databases, which keeps costs roughly linear with growth and lets us isolate our largest merchants so one flash sale can't slow down everyone else. It's about two quarters of work for a four-person team, plus ongoing ownership by the data platform group. The alternative — a managed distributed database — is faster to adopt but roughly 1.5–2× the infrastructure bill and adds latency to checkout. I recommend we build on the proven open-source approach now, and make it the standard every team uses, so we pay this cost once instead of five times."

Principal Interview Signals#

SignalWhat It Sounds Like
Prices the decision"Hand-sharding is ~$150K of tooling and 4 FTE ongoing; distributed SQL is ~$1M/year more infra at our scale. Break-even is around 40 clusters."
Names the one-way doors"The key and the ID format are permanent; host counts and routers are not. I'll spend my review time on the first two."
Designs the org failure posture"Cells by tenant tier limit blast radius to 1/32 of customers; game days prove it quarterly."
Writes the standard"One partition axis, one router, one migration tool — with an exception process, not a mandate without escape."
Knows when not to standardize"Analytics and event logs don't go on this platform — they have different access patterns."

Staff answers that L7 interviewers find insufficient:

  • A flawless resharding workflow for one database, with no view on the other four teams about to build their own.
  • "We'll use Vitess" without the headcount, the migration timeline or the exit cost if Vitess stops fitting.
  • A hot-shard fix that's purely technical, ignoring that whale isolation should be a priced product tier.

Appendices

Appendix A: Partitioning Mechanics in Depth

A.1 Hash mod N — Why It's a Trap#

shard = hash(key) % N
N: 8 → 9   ⇒ a key stays put only if hash % 8 == hash % 9  ⇒ ~1/9 of keys stay; ~89% move

Every resize is a full migration. Acceptable only if N never changes.

A.2 Consistent Hashing#

Each node owns arcs of a ring; with ~100–256 virtual nodes per physical node, load variance drops to a few percent. Adding one node to N moves ~1/(N+1) of keys. Right for caches and leaderless KV stores; weak for relational OLTP because you can't pin entities, range scans are impossible, and the "move" is still an unmanaged data copy.

A.3 Logical Shards + Directory (Default)#

logical  = murmur3(merchant_id) % 4096
physical = override.get(merchant_id) or map[logical]   # map: 4096 entries, ~100 KB
route(query) → physical cluster, stamped with map_version

Moves operate on whole logical shards. Overrides pin individual whales (keep < 1,000 entries).

A.4 Range Partitioning#

Good for time-ordered data and range scans; splits are cheap (split a range at a midpoint). Monotonic keys create a hot tail — mitigate with (time_bucket, hash(entity) % K). Bigtable, HBase and Spanner-family systems auto-split ranges by size and load.

A.5 Entity-Group / Hierarchical Keys#

Child rows carry the parent's key (merchant_id on orders, order_items, inventory) so the whole entity group co-locates. This is the pattern behind Spanner's interleaved tables and every successful hand-sharded schema.

Appendix B: Keys, IDs and Data Model

B.1 Shard-Aware ID Format#

64-bit id = [41 bits ms since epoch][12 bits logical shard][11 bits per-ms sequence]
           ≈ 69 years          4096 shards        2048 ids/ms/shard

Any service can route GetOrder(id) without a lookup; IDs sort by time.

B.2 Denormalize the Partition Axis#

Every table in the entity group stores merchant_id even if derivable. Without it, you can't route child-row queries.

B.3 Reference Data#

Small, read-mostly tables (currencies, categories) are replicated to every shard, updated via a broadcast job. Never join to a "global" shard at request time.

B.4 Global Uniqueness — Claim Tables#

email_claims(email PK, user_id, created_at) sharded by hash(email); claim first, then create the user; sweeper removes orphan claims after 10 minutes.

Appendix C: Cross-Shard Coordination Mechanisms

C.1 Transactional Outbox#

BEGIN;
  INSERT INTO orders (...) VALUES (...);
  INSERT INTO outbox (id, topic, payload) VALUES (uuid, 'order.created', ...);
COMMIT;
-- relay (or CDC on the outbox table) publishes and marks sent; consumers dedupe by id

C.2 Saga with Compensation#

Each step is idempotent keyed by a business ID; each has a compensating action; a reconciliation job closes gaps.

C.3 2PC / XA#

Prepare on all participants, coordinator logs decision, commit. Blocking on coordinator loss after prepare; lock hold ≥ 2 RTTs.

C.4 Quick Comparison#

MechanismAtomicityAvailability CouplingLatencyOperational Burden
Design awayN/ANoneLowestData-modeling effort
Outbox + sagaEventualNone+async secondsCompensation logic
2PC / XAAtomicAll participants + coordinator+5–20 msIn-doubt txn recovery
Distributed SQLSerializableQuorum per range+2–10 ms in-regionVendor/engine expertise
Appendix D: Router Contract and Client Behavior

D.1 Error Semantics#

ErrorMeaningClient Action
WRONG_SHARD(owner, version)Stale mapRefresh map, retry once
SHARD_FROZEN(retry_after_ms)Cutover in progressRetry after 100–500 ms with jitter; give up after 5 s
SHARD_UNAVAILABLEPrimary failoverRetry with backoff up to 30 s; writes are idempotent by request ID
FANOUT_BUDGET_EXCEEDEDScatter-gather quota hitUse derived store

D.2 Idempotency#

Every write carries a client request ID stored with the row, so retries across failover and cutover don't duplicate.

Appendix E: Observability

E.1 Core Metrics#

Per shard:   shard.cpu_pct, shard.qps{op}, shard.p99_ms{op}, shard.size_gb, shard.disk_pct,
             shard.replication_lag_s, shard.top_keys_qps
Router:      router.fanout_ratio, router.map_version, router.map_version_skew, router.wrong_shard_retries
Migration:   migration.copy_mb_s, migration.catchup_lag_s, migration.verify_mismatch_chunks, migration.freeze_ms
Derived:     gsi.lag_seconds, outbox.unpublished_rows, saga.pending_age_seconds

E.2 Critical Alerts#

AlertThresholdSeverity
Shard imbalancemax/median shard.cpu_pct > 2.5 for 30 minTicket
Hot shardshard.cpu_pct > 85% for 5 minPage
Map skewrouter.map_version_skew > 0 for 60 sPage
Freeze stuckmigration.freeze_ms > 10,000Page
GSI laggsi.lag_seconds p99 > 10 for 5 minPage data platform
Diskshard.disk_pct > 85%Page

E.3 Control Plane vs Data Plane#

Data plane (routers, clusters) must survive a control-plane (directory, orchestrator) outage indefinitely with cached state. Control plane outages block change, never traffic.

E.4 Debugging Imbalance#

Rank logical shards by QPS and size; if the top logical shard is > 5× median, look for a whale key; if many logical shards on one cluster are warm, rebalance placement.

Appendix F: Scale Evolution

F.1 What Works at Each Scale#

ScaleApproach
< 2 TB, < 10K w/sSingle primary, replicas, table partitioning
2–20 TB, 10–100K w/sLogical shards on 8–64 clusters, proxy router
20–500 TB, 100K–1M w/sPartitioning platform, automated balancing, cells
Multi-regionHome-region placement, region-aware directory

F.2 What You Don't Build on Day One#

  • Automated shard balancing — manual moves with good tooling are enough until ~32 clusters
  • Cross-region active-active — home-region is enough until residency or regional SLOs demand more
  • A distributed query engine — use the warehouse
  • Synchronous GSIs — never
Appendix G: Multi-Tenancy, Fairness and Cost

G.1 Tenant Tiers#

TierPlacementIsolationPrice Signal
Long tailHash into shared logical shardsPer-tenant write budgetIncluded
GrowthShared clusters, monitoredPer-tenant limits + alertingIncluded
EnterpriseDedicated cluster via overridePhysicalPriced into contract
RegulatedDedicated cluster in required regionPhysical + residencyPremium

G.2 Showback#

Attribute cluster cost to logical shards by size and QPS, and logical shards to tenants and owning teams. Publish monthly; it changes behavior faster than any standard.

  1. Loading the index…