Why This Matters#
Sharding is not a scaling technique you add when the database gets big. It is a permanent decision about which operations stay cheap and which become distributed systems problems. The partition key decides which queries hit one node and which fan out to all of them, which transactions are local and which need coordination, and which customer becomes the hot spot that pages you on Black Friday. Choosing it is one of the few genuinely one-way doors in system design.
Most candidates say "we'll shard by user ID" as a reflex, then keep designing as if there were one database. Staff candidates treat the partition key as the moment the design forks: they check every top access pattern against it, name the queries that become scatter-gather, name the invariants that now span partitions, and say what happens when one key receives 1,000× the average traffic.
This page is the conceptual foundation: strategies, key selection, rebalancing, hot partitions, and cross-partition operations. For the full interview-grade design — resharding a live system, routing tiers, and migration runbooks — see the Database Sharding case study.
The 60-Second Version#
- Don't shard until you must. A single modern primary handles ~1–10 TB and ~10–50K writes/s. Read replicas, caching, and index hygiene defer sharding by years. Sharding costs 2–4 engineers for 2–3 quarters and taxes every future feature.
- Three strategies: range, hash, directory. Range preserves order (great scans, hot tails). Hash spreads load (no range scans). Directory maps keys to shards explicitly (maximum flexibility, one more critical system).
- The partition key must match the dominant access pattern and the dominant transaction boundary. If 90% of queries and all invariants are per-tenant, shard by tenant.
- Rebalancing is the hard part, not initial placement. Use many more logical partitions than physical nodes (e.g., 4,096 logical → 16 nodes) so growth moves whole partitions instead of rehashing keys.
hash(key) % Nis the classic mistake — changing N moves ~all keys. - Hot partitions are inevitable in skewed data. A single partition typically sustains ~1–10K writes/s depending on the engine (DynamoDB documents ~1,000 WCU per partition). The top 0.01% of keys will exceed it. Plan key splitting, write sharding, and caching for them.
- Cross-partition operations are the tax. Scatter-gather reads (latency = slowest shard), cross-shard transactions (2PC or sagas), and global secondary indexes (async or expensive). Minimize them by design, not by optimization.
- Resharding a live system is a migration, not a config change. Dual-write, backfill, verify, cut over — with the same rigor as a schema migration.
How Sharding Works#
The Basic Idea#
Partitioning splits a dataset into disjoint subsets; sharding places those subsets on different nodes. (Many teams use the words interchangeably; the distinction that matters is logical partitions vs physical nodes.) Each record lives in exactly one partition, determined by a function of its partition key. A router — in the client library, a proxy, or the database itself — maps key → partition → node.
The payoff: write throughput, storage, and working-set memory scale roughly linearly with node count, for operations that touch one partition. Everything else gets harder.
Key Terms#
| Term | Meaning | Why It Matters |
|---|---|---|
| Partition key | The attribute(s) that determine placement | Decides locality of every query and transaction |
| Logical partition / vnode / tablet | A unit of data movable as a whole | Rebalancing moves these, not individual keys |
| Physical shard / node | The machine or cluster holding partitions | Capacity unit; failure unit |
| Router / directory | The mapping from key or partition to node | Must be highly available and consistent; stale routes send writes to the wrong place |
| Hot partition / hot key | A partition receiving disproportionate load | Throughput is capped at one node's capacity, regardless of cluster size |
| Scatter-gather | Query sent to all partitions and merged | Cost × N; latency = slowest partition |
| Co-location | Related data sharing a partition key | Enables local joins and transactions |
| Resharding | Changing the number or layout of partitions | Moves data under live traffic |
Where Partitioning Lives#
| Layer | Example | Who Chooses the Key | Rebalancing |
|---|---|---|---|
| Application-level sharding | Many shards of MySQL/Postgres with a routing library | Application team | Manual, tooling-heavy |
| Proxy / middleware | Vitess (MySQL), Citus (Postgres) | Application team via schema/VSchema | Tool-assisted (resharding workflows) |
| Native distributed DB | Cassandra, DynamoDB, CockroachDB, Spanner, MongoDB | Application picks key; DB places and splits | Automatic (splits/vnodes), but key choice still decides hot spots |
| Streams | Kafka topic partitions | Producer's message key | Adding partitions breaks key→partition mapping for existing keys |
| Caches | Redis Cluster (16,384 slots), Memcached client hashing | Client | Slot migration / consistent hashing |
🎯 Staff Insight: "Native distributed databases automate placement, not key choice. DynamoDB will split a hot partition by range; it cannot split a single hot key. The key design is still mine."
Core Strategies#
Strategy 1: Range Partitioning#
Contiguous key ranges map to partitions: [a–f) → P1, [f–m) → P2, ... Used by HBase, Bigtable, Spanner, CockroachDB, and MongoDB ranged sharding.
route(key):
return partition_map.floor_entry(key).partition # binary search on range boundaries
split(partition): # automatic in Bigtable/Spanner/CockroachDB
if partition.size > 512 MB or partition.qps > threshold:
mid = partition.median_key()
create [start, mid) and [mid, end); move one half to a less loaded node
When to use: range scans and ordered access dominate (time ranges per entity, lexicographic scans, "all rows for tenant X"). Failure mode: monotonic keys create a moving hot spot. If the key is a timestamp or auto-increment ID, every new write lands in the last range; one node takes 100% of inserts while the rest idle. Fix: prefix with a high-cardinality attribute ((tenant_id, ts)), or salt/hash the leading component.
Strategy 2: Hash Partitioning#
Partition = hash(key) mapped onto a fixed space of logical partitions or a hash ring. Used by Cassandra (Murmur3 token ring), DynamoDB (partition key hashing), Redis Cluster (CRC16 → 16,384 slots), Kafka (key hash → partition).
NUM_LOGICAL = 4096 # fixed forever; >> number of nodes
route(key):
p = murmur3(key) % NUM_LOGICAL # stable: never depends on node count
return partition_to_node[p] # small table; changes on rebalance
# WRONG: node = hash(key) % num_nodes → going 16 → 17 nodes moves ~94% of keys
When to use: point lookups and writes spread across many keys; no need for range scans on the partition key. Failure mode: range queries become scatter-gather, and hashing does nothing for a single hot key — hash("taylor_swift") always lands on the same partition. See Consistent Hashing for the ring variant and minimal-movement rebalancing.
Strategy 3: Directory (Lookup) Partitioning#
An explicit mapping service records which shard holds each key or key group: tenant_42 → shard_7.
route(tenant_id):
shard = directory_cache.get(tenant_id) # local cache, TTL ~30-60s + invalidation
if shard is null:
shard = directory_service.lookup(tenant_id) # strongly consistent store (etcd/Spanner/Postgres)
return shard
move_tenant(tenant_id, to_shard): # per-tenant migration becomes possible
copy → catch up via CDC → brief write freeze (~seconds) → flip directory → unfreeze
When to use: multi-tenant SaaS with very uneven tenant sizes, data residency requirements, or dedicated shards for top customers. Slack, Notion-style workspace sharding and most B2B platforms end up here. Failure mode: the directory is now a tier-0 dependency. If it's unavailable or stale, every request fails or misroutes. Cache it aggressively, make it strongly consistent, and version it so stale routes are detected (shard rejects writes for tenants it no longer owns).
Strategy Comparison#
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Range | Ordered scans; automatic splitting; locality for prefixes | Sequential keys hot-spot the tail range | On-call during write bursts; app team must design key prefixes |
| Hash | Even spread; simple routing; no hot tail | Range scans scatter; single hot keys unsolved; changing partition count is painful | Readers doing range queries; the team that picked too few partitions |
| Directory | Per-tenant placement, isolation, residency, dedicated shards | Directory is critical infra; extra hop or cache; consistency of routes | Platform team owns directory availability; every request depends on it |
| Hybrid (directory of tenants → hash within) | Tenant isolation plus spread for huge tenants | Two layers of routing to operate | Platform team, in complexity |
🎯 Staff Move: "For B2B, I'd shard by tenant through a directory: 99% of queries and every invariant are tenant-scoped, and the directory lets us move our largest customer to a dedicated shard without touching anyone else. Inside a giant tenant, I'd hash by a secondary key. For a consumer app with user-scoped data and no tenants, hash by
user_idinto 4,096 logical partitions."
Choosing the Partition Key#
This is the hard sub-problem. Everything else — rebalancing, hot spots, cross-partition cost — is downstream of it.
The Four Tests#
| Test | Question | Failing Looks Like |
|---|---|---|
| 1. Access locality | Do the top queries by volume include the key? | 40% of queries become scatter-gather |
| 2. Transaction locality | Do invariants (uniqueness, balances, counters) live within one key's data? | Every checkout is a cross-shard transaction |
| 3. Cardinality & distribution | Are there enough distinct values, spread evenly? | country as a key: 200 values, one of them is 40% of traffic |
| 4. Growth & skew | Can any single key outgrow one partition? | One tenant grows to 30% of total data |
# Worked example: messaging app
candidates:
user_id → inbox reads local ✓; a message to a group writes N inboxes ✗ (fan-out on write)
conversation_id → reading a conversation local ✓; inbox (list my conversations) is scatter ✗
(conversation_id, time_bucket) → bounded partitions ✓ for very long chats
decision:
messages partitioned by (conversation_id, month) -- primary read: "load this chat"
user_inbox partitioned by user_id -- denormalized projection: "my chats"
=> two partition keys, two tables, one async projection; each hot path is single-partition
The pattern generalizes: when two top access patterns want different keys, keep two copies with two keys and maintain one from the other. That is a data-modeling decision (see Data Modeling & Schema Design), and it's usually cheaper than making either path scatter.
Compound Keys and Bucketing#
- Compound partition key:
(tenant_id, entity_type)or(device_id, day)bounds partition size for unbounded data. - Partition key + sort key: DynamoDB/Cassandra split placement (partition key) from order within the partition (sort/clustering key):
PK=user_id, SK=created_atgives single-partition time-range queries per user. - Time bucketing: for append-only series,
(entity, hour|day|month)keeps partitions under engine limits (e.g., Cassandra guidance of ~100 MB) and makes retention a partition drop.
Rebalancing#
Data grows, nodes are added, and hot spots move. Rebalancing moves partitions under live traffic.
Fixed Logical Partitions (The Default)#
Create far more logical partitions than nodes — e.g., 256–4,096 — at day one. Nodes own sets of partitions. Adding a node steals whole partitions from others; key → partition never changes.
| Approach | Keys Moved When Adding 1 Node to 16 | Operational Model | Used By |
|---|---|---|---|
hash % N | ~94% | Full reshuffle | Nobody, deliberately |
| Consistent hashing with vnodes | ~1/17 ≈ 6% | Tokens reassigned | Cassandra, Dynamo-style systems |
| Fixed logical partitions | ~1/17 ≈ 6% (whole partitions) | Move partitions | Redis Cluster slots, Elasticsearch shards (fixed at index creation), Couchbase vBuckets |
| Dynamic range splits | Only the split partition | Automatic split + move | Bigtable, HBase, Spanner, CockroachDB, DynamoDB |
The catch with fixed partitions: the count is chosen once. Elasticsearch's primary shard count is fixed at index creation (changing it requires split/shrink or reindex). Kafka lets you add partitions, but existing keys then map to different partitions, breaking per-key ordering. Choose the count for 3–5 years of growth.
The Mechanics of Moving a Partition#
move(partition p, from A, to B):
1. snapshot p on A → stream to B # bulk copy, throttled (e.g., 50-100 MB/s)
2. tail A's change log for p → apply on B # catch up until lag < ~1s
3. briefly block writes to p on A (ms-seconds) # or use a lease/epoch handoff
4. drain final changes; bump routing epoch # p now owned by B at epoch e+1
5. A rejects writes for p with epoch <= e # stale routers get a redirect, refresh
6. clean up p on A after a grace period
Rules: throttle movement (rebalancing competes with production I/O); move during low traffic; never move many partitions at once after a node failure (a "rebalance storm" can saturate the network and cause the next failure); use epochs/fencing so a stale router can't write to the old owner.
Automatic vs Manual Rebalancing#
Automatic rebalancing sounds strictly better. It isn't: an automatic rebalancer reacting to a slow (not dead) node can move terabytes, amplify load, and turn a partial failure into a cluster-wide one. Common production posture: automatic splits for growth, human-approved (or rate-limited) moves for failure recovery.
Hot Partitions#
Skew is the natural state of real data: Zipfian popularity, whale tenants, celebrity accounts, viral items, the current time bucket.
| Hot-Spot Type | Example | Mitigation | Cost |
|---|---|---|---|
| Hot write key | Global counter, trending post's like count | Write sharding: key#0..key#N-1, sum on read; or aggregate in memory and flush | Reads fan out to N sub-keys; eventual totals |
| Hot read key | Celebrity profile | Cache (L1 + L2), read replicas of the partition, request coalescing | Staleness |
| Monotonic key | Timestamp or sequence as leading key | Prefix with high-cardinality attribute; hash/salt the prefix | Range scans across the salt become scatter |
| Whale tenant | One customer = 30% of data | Directory: move to dedicated shard; hash within tenant | Dedicated infra cost; per-tenant ops |
| Temporal hot bucket | "Today" partition in time series | Sub-bucket by (entity, hour) or add hash suffix | More partitions to query for ranges |
# Write sharding a hot counter (N sub-keys)
increment(post_id):
shard = random(0, N-1) # N = 10-100 for a viral post
db.increment(f"likes:{post_id}#{shard}", 1)
read_likes(post_id):
return sum(db.get(f"likes:{post_id}#{i}") for i in 0..N-1) # or cache the sum for 1-5s
# Adaptive: start N=1; when the key's write rate exceeds ~50% of partition capacity,
# raise N and record it in metadata so readers know how many sub-keys to sum.
🎯 Staff Move: "Hashing spreads keys, not load. I'd assume the top 0.01% of keys will exceed a partition's ~1K writes/sec and plan for it: detect with per-key metrics, write-shard the counter, cache the reads, and give whale tenants their own shard via the directory."
Cross-Partition Operations#
Every operation that touches more than one partition pays a tax. Design to make them rare, then make them explicit.
| Operation | Mechanism | Cost | Staff Default |
|---|---|---|---|
| Query without the partition key | Scatter-gather to all N | N× work; p99 = slowest shard; one slow shard stalls all | Add a projection keyed for that query, or a search index via CDC |
| Secondary index lookup | Local index (scatter) or global index (async or 2PC) | See Database Indexing | Global mapping table for uniqueness; search for everything else |
| Join across partitions | Application-side join, broadcast small tables | Network round trips; memory | Co-locate by shared key (tenant_id); replicate small reference tables to every shard |
| Multi-partition transaction | 2PC (Spanner, CockroachDB), or sagas with compensation | 2PC: +1 round trip and coordinator risk; sagas: temporary inconsistency | Choose the key so invariants are single-partition; sagas for the rest |
| Global aggregates | Scatter-gather or pre-aggregation | Expensive at query time | Stream aggregates into a separate store |
| Global uniqueness | Dedicated partition keyed by the unique value | Extra write, ordering | Claim-then-create with a conditional put |
# Transfer between two accounts on different shards: saga with idempotent steps
transfer(tx_id, from, to, amount):
shard(from).debit(tx_id, from, amount) # idempotent on tx_id; records PENDING_OUT
try:
shard(to).credit(tx_id, to, amount) # idempotent on tx_id
shard(from).mark_complete(tx_id)
except PermanentFailure:
shard(from).refund(tx_id) # compensating action
# Invariant "money is conserved" holds eventually; a reconciler scans PENDING_OUT older than 60s.
🎯 Staff Insight: "If more than ~5–10% of our transactions are cross-shard, the partition key is wrong. I'd rather re-key than build a faster 2PC."
Visual Guide#
Should You Shard, and How?#
Request Routing with a Directory and Epochs#
Logical Partitions Over Physical Nodes#
Implementation Patterns#
Start with Logical Sharding on One Node#
Before you need physical shards, make the code shard-aware: every table carries the partition key, every query includes it, IDs are globally unique (not per-database auto-increment), and a routing function exists — even though it returns the same database for every key. When the day comes, moving logical partitions to new nodes is an operational task instead of a rewrite. Instagram's early engineering blog described exactly this: thousands of logical shards mapped onto a handful of physical Postgres servers, with IDs that embed the logical shard.
Shard-Aware IDs#
Embed the partition in the ID so any service can route without a lookup:
id (64 bits) = [41 bits ms since epoch][13 bits logical shard][10 bits sequence]
route(id) = shard_bits(id) # no directory call on the read path
Time-ordered for B-tree locality, routable without a lookup, globally unique. The tradeoff: an entity can't move to another logical shard without changing its ID — so the logical shard count must be large enough to rebalance at the logical level.
Co-location and Reference Data#
Give related tables the same partition key so joins and transactions stay local (Citus calls these co-located distributed tables; Vitess uses keyspace IDs). Small, rarely changing tables (countries, plans, feature definitions) are replicated to every shard so joins against them never fan out.
Fan-out Guardrails#
When scatter-gather is unavoidable: issue in parallel with a per-shard timeout; cap concurrency; return partial results with a partial=true flag for non-critical reads; hedge the slowest shard; and push aggregation down to shards so you merge small results, not raw rows.
Resharding Playbook (Summary)#
Resharding a live system follows expand → migrate → contract: provision new shards, dual-write or CDC-replicate, backfill, verify with checksums per partition, switch reads, switch writes with a brief freeze or epoch flip, then decommission. The full procedure, including failure handling and rollback, is in the Database Sharding case study.
The Numbers in Context#
| Number | Value | What It Means for Your Design |
|---|---|---|
| Single primary ceiling | ~1–10 TB, ~10–50K writes/s | Most systems should not shard; replicas and caching come first |
| Per-partition throughput | ~1–10K writes/s depending on engine; DynamoDB documents 1,000 WCU / 3,000 RCU per partition | One hot key caps at this regardless of cluster size |
| Logical partitions | 256–4,096 (Redis Cluster: 16,384 slots) | Enough to rebalance for 3–5 years without changing key → partition |
| Partitions per node | ~10–100+ | Finer granularity = smoother rebalancing, more metadata |
| Target partition size | ~10–50 GB (range DBs split far smaller, e.g., hundreds of MB) | Big enough to amortize overhead, small enough to move in minutes |
| Move rate | ~50–200 MB/s throttled | A 50 GB partition moves in ~5–15 minutes |
hash % N reshuffle | ~(N)/(N+1) of keys move | 16 → 17 nodes moves ~94% of keys |
| Scatter-gather at 64 shards | 64× requests; p99 ≈ max of 64 | If per-shard p99 = 10ms, fan-out p50 ≈ 10ms |
| Cross-shard transaction share | Aim for < ~5–10% | Above this, re-examine the partition key |
| Whale tenant threshold | > ~5–10% of a shard's capacity | Candidate for dedicated shard via directory |
| Sharding project cost | ~2–4 engineers × 2–3 quarters | Every quarter of deferral is worth real effort |
How This Shows Up in Interviews#
Scenario 1: "How would you scale the database?"#
L5 jumps to "shard by user ID." Staff first sizes: "At 8K writes/s and 2 TB, one primary with replicas holds for ~18 months at current growth. I'd design logically sharded now — partition key on every table, globally unique IDs, routing function — and physically shard when we're within 12 months of the ceiling." Then name the key and check it against the top queries and invariants.
Scenario 2: "One shard is at 95% CPU while the others are at 30%" (Full Walkthrough)#
Step 1 — Is it a hot partition or a hot key? "I'd break down load on that shard by logical partition, then by key. If one logical partition dominates, it's a placement problem — move it or split it. If one key dominates, moving it just moves the fire."
Step 2 — Hot partition (placement). "Several busy logical partitions happened to land on one node. Move 2–3 of them to the coolest nodes, throttled at ~100 MB/s, during off-peak, with epoch-fenced handoff. Then fix the placement policy to balance on load, not just partition count."
Step 3 — Hot key (skew). "If it's one tenant or one key, the fix depends on what it is. A whale tenant: move it to a dedicated shard via the directory, and hash within the tenant if it outgrows one node. A hot counter: write-shard it into 10–100 sub-keys and cache the sum for a few seconds. A hot read: L1 cache with a short TTL plus request coalescing."
Step 4 — Protect the neighbors now. "While we fix it, per-tenant rate limits on that shard so the whale's traffic can't starve the other ~200 tenants co-located with it — their SLO matters as much as the whale's."
Step 5 — Detect earlier next time. "Metrics: shard.cpu_utilization spread (max/median > 2 alerts), partition.qps top-K, key.write_rate top-K sampled. The storage platform owns placement and the rebalancer; the product team owns key design and the whale-tenant policy with account management."
Why this is a Staff answer: it separates placement from skew, applies the right fix to each, protects co-tenants during mitigation, and assigns detection and ownership.
Scenario 3: "Now we need a query that doesn't include the partition key"#
Name the cost (scatter to all N shards), then decide by QPS and freshness: < ~10 QPS and tolerant of seconds of latency → scatter with timeouts and partial results; high QPS → a projection or global index keyed for the new query, maintained via CDC; free-form search → search index. Mention that the answer changes if the query needs strong consistency (e.g., uniqueness), which requires a synchronous mapping.
Scenario 4: "We picked the wrong partition key"#
It happens. The fix is a resharding migration: new cluster with the new key, CDC from old to new, backfill, verify, dual-read comparison, cutover by traffic slice. Cost: quarters. Point to the Database Sharding case study and say what you'd do differently: validate the key against the top 10 access patterns and a 3-year growth model before committing.
What Interviewers Probe#
| After You Say... | They Will Ask... | (What They're Evaluating) |
|---|---|---|
| "Shard by user ID" | "Now show me a user's group conversations and the group's members." | Whether you check the key against every top access pattern |
| "Consistent hashing handles rebalancing" | "One key gets 50K writes/sec. Where does it go?" | Keys vs load; single-key hot spots |
| "We'll add shards as we grow" | "Walk me through going from 16 to 24 shards with live traffic." | Logical partitions, throttled moves, epoch fencing |
| "Use a distributed transaction" | "What fraction of your writes are cross-shard, and what's the p99 cost?" | Designing invariants to be partition-local |
| "A directory maps tenants to shards" | "The directory is down. What happens?" | Treating routing as tier-0 infrastructure |
| "Timestamp-prefixed keys for range scans" | "Where does every new write land?" | Monotonic-key hot tails in range partitioning |
Advanced Patterns#
| Pattern | How It Works | When to Use |
|---|---|---|
| Cell-based architecture | Entire stacks (app + DB) per cell; tenants assigned to cells | Blast-radius isolation at large scale; AWS publicly advocates it |
| Hierarchical keys | (tenant, entity) with directory at tenant level, hash below | B2B with whale tenants |
| Adaptive write sharding | Increase sub-key count for keys crossing a rate threshold | Viral counters, trending items |
| Geo-partitioning | Partition by region for residency and latency | Regulatory requirements; region-local users |
| Reference table replication | Small tables copied to every shard | Joins against lookup data |
| Shard-aware IDs | Partition encoded in the ID | Lookup-free routing |
| Split-and-merge | Automatic range splitting on size/load; merge when cold | Range-partitioned distributed SQL |
| Shadow reads during reshard | Read from old and new, compare, serve old | Verifying a migration before cutover |
Failure Modes & Operational Reality#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Hot partition / hot key | partition.qps top-K; throttling errors (e.g., DynamoDB ProvisionedThroughputExceeded); shard CPU spread | Every tenant co-located on that shard | Split/move partition; write-shard key; cache; dedicated shard | Product team (key) + storage platform (placement) |
| Stale routing | router.moved_redirects ↑; writes to wrong shard | Tenants being moved | Epoch fencing; shard rejects non-owned keys; route cache invalidation | Storage platform |
| Directory outage | directory.lookup_errors; cache miss rate ↑ | All requests needing uncached routes | Long-lived local caches, stale-OK reads, multi-region directory | Storage platform |
| Rebalance storm | rebalance.bytes_in_flight ↑, network saturation, p99 ↑ cluster-wide | Entire cluster | Rate-limit moves; manual approval after failures | Storage on-call |
| Scatter-gather tail latency | Fan-out query p99 ≫ per-shard p99 | Any endpoint using the query | Projection keyed for the query; partial results; hedging | Owning service team |
| Cross-shard saga stuck | saga.pending_age_seconds p99 ↑ | Money/inventory in limbo | Reconciler, compensation, alerts at > 60s | Owning service team |
Failure Scenario: The Whale Tenant Onboarding#
t=0 Sales closes a customer 40× larger than the median tenant. Hash placement puts them on shard 11.
t=+2d Their bulk import starts: 30K writes/s against shard 11, sized for ~8K.
t=+2d+5m Shard 11 p99 → 2s; the 180 other tenants on shard 11 see timeouts. Pages fire.
t=+2d+20m On-call throttles the import; co-tenants recover. The whale's import stalls.
t=+2d+3h Emergency: provision a dedicated shard, migrate the tenant via CDC, flip directory entry.
t=+3d Import resumes on dedicated shard at full speed.
Detection: alert when a single tenant exceeds ~20% of a shard's capacity; pre-onboarding size estimates from sales. Prevention: a whale-tenant policy — contracts above a size threshold trigger placement on a dedicated or lightly loaded shard before the import; per-tenant write rate limits on shared shards. Owner: storage platform owns placement; account management owns flagging large deals; product owns import throttling.
The Principal Lens#
Why L7 Sees This Problem Differently#
At Staff level, sharding is about choosing a good key and rebalancing safely for one system. At Principal level, sharding is a company-wide architecture that determines blast radius, residency, tenant isolation, and unit economics — and it's usually being solved five different ways by five teams. The Principal asks whether the org should have one partitioning model (tenant → cell) that every stateful service follows, so a customer lives in the same cell across all services, a cell failure affects a bounded fraction of customers, and "move this customer to the EU" is one operation instead of fifteen migrations.
The Org-Level Fault Line#
Per-service sharding vs a shared cell architecture. Letting each team shard its own store by its own key maximizes local fit. A shared cell model (tenants assigned to cells; each cell contains every service's partition for those tenants) aligns blast radius and residency but constrains every team's key choice and requires a platform to operate cells.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Each service shards independently | Best local key per service | A customer spans different shards everywhere; no consistent blast radius; residency = N migrations | Every team re-solves routing and rebalancing; SRE can't reason about correlated failures |
| Shared cell architecture | Bounded blast radius per cell; one tenant-placement system; residency per cell | Cross-cell features (global search, analytics) need separate pipelines; cell platform is heavy | Platform org (~5–15 FTE); product teams adopt tenant-scoped design |
| Shared routing/directory library, per-service stores | Consistent tenant placement primitives without full cells | Partial alignment only | Platform (~2–4 FTE) |
Cost Model#
Assumptions: engineer ≈ $25K/month; managed database nodes ~$2–8K/month each at the sizes shown; figures are order-of-magnitude.
| Scale | Topology | Infra / Month | People | On-Call Load |
|---|---|---|---|---|
| Pre-shard (2 TB, 8K writes/s) | 1 primary + 2 replicas | ~$10–20K | 0.5 FTE DB ops | Low; failover drills quarterly |
| Sharded (40 TB, 150K writes/s) | 16–32 shards × 3 replicas + directory | ~$150–350K | Sharding project: 3 FTE × 3 quarters (~$700K one-time); 2–3 FTE ongoing | Rebalancing, hot-shard pages ~2–4/month |
| Cells (1 PB, multi-region, 5M tenants) | 50–200 cells, each with full stack | ~$2–6M | Cell platform 8–15 FTE (~$200–375K/month) | Per-cell on-call automation; cell-level game days |
The lever: every quarter sharding is deferred through replicas, caching, and index hygiene saves the project's ongoing tax. But once you shard, designing it as tenant → cell from the start avoids a second, larger re-architecture later.
The 3-Year Evolution Path#
One-Way Doors vs Two-Way Doors#
| Decision | Reversibility | Cost to Reverse |
|---|---|---|
| Partition key | One-way | Full reshard; quarters of migration |
| Number of fixed logical partitions | One-way-ish | Reindex / rehash all data (Elasticsearch shards, Kafka partitions) |
| ID format (shard-aware vs opaque) | One-way | IDs are stored and exposed everywhere |
| Physical node count | Two-way | Move partitions |
| Directory vs pure hash routing | One-way-ish | Adding a directory later requires backfilling mappings and a routing migration |
| Moving one tenant to a dedicated shard | Two-way | Move it back |
| Adopting a cell architecture | One-way | Org-wide operating model |
The Standard I'd Write#
RFC-STORE-007: Partitioning Standard for Stateful Services
Scope: All stateful services storing customer data, new and existing (existing services comply at next major storage change).
Mandatory (MUST):
- Tenant-scoped data MUST carry
tenant_idas the leading partition key component.- Services MUST route through the shared placement directory; no service-private tenant → shard mappings.
- Hash-partitioned stores MUST use a fixed logical partition count ≥ 64× the initial node count;
hash % nodesrouting is prohibited.- Partition ownership changes MUST be epoch-fenced; shards MUST reject writes for partitions they don't own.
- Per-tenant rate limits MUST exist on shared shards.
Recommended (SHOULD): shard-aware IDs; reference tables replicated per shard; cross-shard operations < 10% of transactions, reviewed if higher.
Exceptions: global, non-tenant data (catalogs, public content) may use other keys with storage platform review.
Success metrics: % of customer data placed via directory (target 90% in 4 quarters); hot-shard Sev2+ incidents per quarter (target ≤ 1); time to relocate a tenant (target < 1 day).
What I'd Tell the VP#
"Our database will hit its limits in about a year, and splitting it is a significant project — roughly three engineers for most of a year. I want to do it once, in a way that also solves problems we know are coming: large customers who slow down everyone else, European customers who need their data kept in the EU, and outages that currently affect all customers at once. The plan is to place each customer in a defined slice of our infrastructure, so we can move big or regulated customers independently and limit how many customers any single failure touches. This costs more upfront than a minimal split, but it avoids a second re-architecture in two to three years."
Principal Interview Signals#
| Signal | What It Sounds Like |
|---|---|
| Connects sharding to blast radius | "A partition is also a failure domain. I want a shard outage to hit 1/32 of customers, not a random slice of every feature." |
| Prices the project and its deferral | "Sharding is ~$700K of engineering plus ongoing tax. Six months of deferral via replicas and index cleanup is worth ~2 engineer-months of effort." |
| Aligns partitioning across services | "If orders are sharded by tenant and billing by account, residency is two migrations. One tenant-placement system for all of them." |
| Designs for the whale | "Sales needs a trigger: deals above a size threshold get placement review before the contract is signed." |
| Names one-way doors | "The partition key and the logical partition count are forever. Node count isn't — I'd over-provision partitions, not nodes." |
Staff answers that L7 interviewers find insufficient:
- "Hash by user ID with consistent hashing" — technically sound, but silent on tenant isolation, residency, and blast radius as business requirements.
- "We'll rebalance when a shard gets hot" — reactive; no whale-tenant policy, no per-tenant limits, no ownership across sales and platform.
- "Each service picks its own shard key" — locally optimal, globally produces misaligned failure domains and N migrations for every residency deal.
🧭 Principal Move: "I'd make tenant placement a platform capability — one directory, one relocation workflow — so that 'move this customer' or 'isolate this customer' is a ticket, not a quarter."
In the Wild#
Instagram: Logical Shards on Postgres#
Instagram's engineering blog described sharding Postgres into thousands of logical shards mapped to far fewer physical servers, with 64-bit IDs that encode creation time and the logical shard ID (generated inside Postgres via a PL/pgSQL function). Moving logical shards between servers let them add capacity without re-keying data.
Staff insight: decoupling logical from physical partitions on day one is what makes rebalancing boring. Cite it when you propose 4,096 logical partitions on 16 nodes.
YouTube / Vitess: Sharding MySQL at Scale#
Vitess was built at YouTube to scale MySQL horizontally, adding a routing layer (VTGate), keyspace IDs for sharding, and resharding workflows that split shards while serving traffic. It was later open-sourced, became a CNCF project, and has been adopted by other large MySQL users (Slack and GitHub have publicly discussed using it).
Staff insight: sharding is a platform, not a library call. The routing layer, resharding workflow, and cross-shard query restrictions are where the engineering goes.
Amazon DynamoDB: Partitions, Adaptive Capacity, and Hot Keys#
DynamoDB documents per-partition throughput limits (1,000 write units and 3,000 read units per second per partition) and provides adaptive capacity and split-for-heat to redistribute load across partitions. Its guidance nonetheless emphasizes high-cardinality partition keys and write sharding for hot keys, because a single partition key value cannot be split.
Staff insight: managed databases automate placement and splitting; they cannot fix a key whose single value is hot. Key design remains the application's responsibility.
Staff Calibration#
What Staff Engineers Say (That Seniors Don't)#
| Concept | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| When to shard | "When the DB gets big, shard it" | "Not until within ~12 months of a measured ceiling; design logically sharded now" | "Sharding is ~$700K plus ongoing tax; I'd fund deferral first and design the eventual split as tenant cells" |
| Key choice | "Shard by user ID" | "Key must satisfy access locality, transaction locality, cardinality, and skew — here's how user ID fares on each" | "One tenant-placement model across all services so failures, residency, and moves align" |
| Rebalancing | "Use consistent hashing" | "4,096 fixed logical partitions, throttled moves with epoch fencing; hash % N is off the table" | "Relocation is a self-service platform workflow with an SLA" |
| Hot spots | "Add more shards" | "Hashing spreads keys, not load; write-shard hot keys, cache hot reads, dedicate shards to whales" | "A whale-tenant policy that starts at the sales contract, not the pager" |
| Cross-shard ops | "Use distributed transactions" | "Choose the key so invariants are local; sagas for the rest; > 10% cross-shard means re-key" | "Cross-cell features get dedicated pipelines; I'd budget them explicitly" |
Why "Key choice" separates levels
"Shard by user ID" is often right. The L6 difference is proving it against the four tests and naming the queries that become scatter-gather and the invariants that become cross-shard. The L7 difference is noticing that fifteen services each picking a locally good key produces fifteen misaligned failure domains.
Why "Hot spots" separates levels
Adding shards helps when load is spread evenly. It does nothing for a single hot key, which is always on one partition. L6 names the specific fix for each hot-spot type. L7 moves detection upstream — to the sales process and tenant onboarding — so the whale never lands on a shared shard unannounced.
Staff Sentence Templates#
"I'd partition by [key] because [top access pattern] and [invariant] are both scoped to it; the cost is that [query] becomes scatter-gather, which is fine at [QPS]."
"We create [N] logical partitions on [M] nodes, so growth moves whole partitions at [MB/s] without changing any key's placement."
"The hottest key will see [rate], above a partition's [capacity], so [write sharding / caching / dedicated shard] handles it, detected by [metric]."
"Cross-shard [operation] is under [percent] of traffic, so I'll use [saga / 2PC / projection] there and keep everything else local."
Common Interview Traps#
- Sharding too early. Proposing 64 shards for 200 GB of data.
hash(key) % N. Changing N moves almost everything.- Monotonic partition keys in range-partitioned stores. All writes hit the last range.
- Ignoring single hot keys. "Consistent hashing handles hot spots" — it doesn't.
- Forgetting cross-shard invariants. Uniqueness, balances, and counters that span keys need explicit design.
- Treating resharding as a config change. It's a data migration with dual-writes, verification, and cutover.
- Per-database auto-increment IDs. They collide across shards.
- No per-tenant limits on shared shards. One tenant's import becomes everyone's outage.
Practice Drill#
Prompt: "You run a B2B analytics product on a single Postgres with 6 TB and 25K writes/s. Your top 3 customers produce 35% of the writes, and a new EU customer requires EU data residency. Design the partitioning approach."
Staff Answer
All queries and invariants are tenant-scoped, and tenant sizes are extremely uneven — so tenant is the partition key and a directory is the router, not a pure hash. Step 1: make the schema tenant-leading everywhere (tenant_id first in every PK and index) and switch to globally unique, time-ordered IDs. Step 2: build a strongly consistent directory (tenant_id → cluster, epoch) with aggressive client caching and epoch-fenced writes. Step 3: place the three whales on dedicated clusters — each is ~10% of writes, so each gets its own primary with room to grow; within a whale, hash on a secondary key if it ever outgrows one node. The long tail goes onto 4–8 shared clusters by hash of tenant_id into 4,096 logical partitions, with per-tenant write limits on shared clusters. Step 4: the EU customer is placed on an EU-region cluster via the same directory; cross-tenant analytics and billing read from a warehouse fed by per-region CDC, with EU data processed in-region. Migration: CDC-replicate each tenant to its new home, verify with per-table checksums, brief write freeze (seconds), flip the directory entry, move whales first. Metrics: shard.write_utilization per cluster, tenant.write_rate top-K, directory.lookup_p99, router.moved_redirects. Storage platform owns the directory and relocation; product owns key design and per-tenant limits.
Why this is L6:
- Picks the key from access and transaction locality, and picks directory routing because of skew.
- Isolates whales and protects the long tail with per-tenant limits.
- Uses one mechanism — the directory — for both load and residency, with epoch fencing and a verifiable migration.
What L7 adds:
- Frames tenant placement as a platform capability every stateful service adopts, so the next residency deal doesn't need a migration per service.
- Prices it:
3 engineers × 2–3 quarters, plus dedicated whale clusters ($30–60K/month), against the EU contract value and reduced incident risk to the long tail. - Proposes the policy with sales and legal: size thresholds that trigger dedicated placement, and residency as a priced product tier.
Where This Appears#
- Database Sharding — the full case study: resharding live systems, routing tiers, migration runbooks, and drills
- Consistent Hashing — the ring, virtual nodes, and minimal-movement rebalancing
- Distributed Caching — sharded caches, hot keys, and slot migration
- Leaderboard — hot counters and write sharding
- Chat Messaging — partitioning by conversation vs user, and time-bucketed partitions
- Scaling Writes — partitioning as the primary write-scaling lever
- Database Indexing — local vs global secondary indexes across partitions
- Data Modeling & Schema Design — keys, aggregates, and projections that keep queries single-partition
Related Technologies: DynamoDB · Cassandra · PostgreSQL · Apache Kafka · Redis