Hiring BarSupport

Sharding & Partitioning

Foundation31 min read4 diagrams

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) % N is 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#

TermMeaningWhy It Matters
Partition keyThe attribute(s) that determine placementDecides locality of every query and transaction
Logical partition / vnode / tabletA unit of data movable as a wholeRebalancing moves these, not individual keys
Physical shard / nodeThe machine or cluster holding partitionsCapacity unit; failure unit
Router / directoryThe mapping from key or partition to nodeMust be highly available and consistent; stale routes send writes to the wrong place
Hot partition / hot keyA partition receiving disproportionate loadThroughput is capped at one node's capacity, regardless of cluster size
Scatter-gatherQuery sent to all partitions and mergedCost × N; latency = slowest partition
Co-locationRelated data sharing a partition keyEnables local joins and transactions
ReshardingChanging the number or layout of partitionsMoves data under live traffic

Where Partitioning Lives#

LayerExampleWho Chooses the KeyRebalancing
Application-level shardingMany shards of MySQL/Postgres with a routing libraryApplication teamManual, tooling-heavy
Proxy / middlewareVitess (MySQL), Citus (Postgres)Application team via schema/VSchemaTool-assisted (resharding workflows)
Native distributed DBCassandra, DynamoDB, CockroachDB, Spanner, MongoDBApplication picks key; DB places and splitsAutomatic (splits/vnodes), but key choice still decides hot spots
StreamsKafka topic partitionsProducer's message keyAdding partitions breaks key→partition mapping for existing keys
CachesRedis Cluster (16,384 slots), Memcached client hashingClientSlot 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#

StrategyWhat WorksWhat BreaksWho Pays
RangeOrdered scans; automatic splitting; locality for prefixesSequential keys hot-spot the tail rangeOn-call during write bursts; app team must design key prefixes
HashEven spread; simple routing; no hot tailRange scans scatter; single hot keys unsolved; changing partition count is painfulReaders doing range queries; the team that picked too few partitions
DirectoryPer-tenant placement, isolation, residency, dedicated shardsDirectory is critical infra; extra hop or cache; consistency of routesPlatform team owns directory availability; every request depends on it
Hybrid (directory of tenants → hash within)Tenant isolation plus spread for huge tenantsTwo layers of routing to operatePlatform 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_id into 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#

TestQuestionFailing Looks Like
1. Access localityDo the top queries by volume include the key?40% of queries become scatter-gather
2. Transaction localityDo invariants (uniqueness, balances, counters) live within one key's data?Every checkout is a cross-shard transaction
3. Cardinality & distributionAre there enough distinct values, spread evenly?country as a key: 200 values, one of them is 40% of traffic
4. Growth & skewCan 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_at gives 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.

ApproachKeys Moved When Adding 1 Node to 16Operational ModelUsed By
hash % N~94%Full reshuffleNobody, deliberately
Consistent hashing with vnodes~1/17 ≈ 6%Tokens reassignedCassandra, Dynamo-style systems
Fixed logical partitions~1/17 ≈ 6% (whole partitions)Move partitionsRedis Cluster slots, Elasticsearch shards (fixed at index creation), Couchbase vBuckets
Dynamic range splitsOnly the split partitionAutomatic split + moveBigtable, 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 TypeExampleMitigationCost
Hot write keyGlobal counter, trending post's like countWrite sharding: key#0..key#N-1, sum on read; or aggregate in memory and flushReads fan out to N sub-keys; eventual totals
Hot read keyCelebrity profileCache (L1 + L2), read replicas of the partition, request coalescingStaleness
Monotonic keyTimestamp or sequence as leading keyPrefix with high-cardinality attribute; hash/salt the prefixRange scans across the salt become scatter
Whale tenantOne customer = 30% of dataDirectory: move to dedicated shard; hash within tenantDedicated infra cost; per-tenant ops
Temporal hot bucket"Today" partition in time seriesSub-bucket by (entity, hour) or add hash suffixMore 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.

OperationMechanismCostStaff Default
Query without the partition keyScatter-gather to all NN× work; p99 = slowest shard; one slow shard stalls allAdd a projection keyed for that query, or a search index via CDC
Secondary index lookupLocal index (scatter) or global index (async or 2PC)See Database IndexingGlobal mapping table for uniqueness; search for everything else
Join across partitionsApplication-side join, broadcast small tablesNetwork round trips; memoryCo-locate by shared key (tenant_id); replicate small reference tables to every shard
Multi-partition transaction2PC (Spanner, CockroachDB), or sagas with compensation2PC: +1 round trip and coordinator risk; sagas: temporary inconsistencyChoose the key so invariants are single-partition; sagas for the rest
Global aggregatesScatter-gather or pre-aggregationExpensive at query timeStream aggregates into a separate store
Global uniquenessDedicated partition keyed by the unique valueExtra write, orderingClaim-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?#

Diagram: Should You Shard, and How?

Request Routing with a Directory and Epochs#

Diagram: Request Routing with a Directory and Epochs

Logical Partitions Over Physical Nodes#

Diagram: 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#

NumberValueWhat It Means for Your Design
Single primary ceiling~1–10 TB, ~10–50K writes/sMost 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 partitionOne hot key caps at this regardless of cluster size
Logical partitions256–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 throttledA 50 GB partition moves in ~5–15 minutes
hash % N reshuffle~(N)/(N+1) of keys move16 → 17 nodes moves ~94% of keys
Scatter-gather at 64 shards64× requests; p99 ≈ max of 64If per-shard p99 = 10ms, fan-out p50 ≈ 10ms
Cross-shard transaction shareAim for < ~5–10%Above this, re-examine the partition key
Whale tenant threshold> ~5–10% of a shard's capacityCandidate for dedicated shard via directory
Sharding project cost~2–4 engineers × 2–3 quartersEvery 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#

PatternHow It WorksWhen to Use
Cell-based architectureEntire stacks (app + DB) per cell; tenants assigned to cellsBlast-radius isolation at large scale; AWS publicly advocates it
Hierarchical keys(tenant, entity) with directory at tenant level, hash belowB2B with whale tenants
Adaptive write shardingIncrease sub-key count for keys crossing a rate thresholdViral counters, trending items
Geo-partitioningPartition by region for residency and latencyRegulatory requirements; region-local users
Reference table replicationSmall tables copied to every shardJoins against lookup data
Shard-aware IDsPartition encoded in the IDLookup-free routing
Split-and-mergeAutomatic range splitting on size/load; merge when coldRange-partitioned distributed SQL
Shadow reads during reshardRead from old and new, compare, serve oldVerifying a migration before cutover

Failure Modes & Operational Reality#

FailureDetection SignalBlast RadiusMitigationOwner
Hot partition / hot keypartition.qps top-K; throttling errors (e.g., DynamoDB ProvisionedThroughputExceeded); shard CPU spreadEvery tenant co-located on that shardSplit/move partition; write-shard key; cache; dedicated shardProduct team (key) + storage platform (placement)
Stale routingrouter.moved_redirects ↑; writes to wrong shardTenants being movedEpoch fencing; shard rejects non-owned keys; route cache invalidationStorage platform
Directory outagedirectory.lookup_errors; cache miss rate ↑All requests needing uncached routesLong-lived local caches, stale-OK reads, multi-region directoryStorage platform
Rebalance stormrebalance.bytes_in_flight ↑, network saturation, p99 ↑ cluster-wideEntire clusterRate-limit moves; manual approval after failuresStorage on-call
Scatter-gather tail latencyFan-out query p99 ≫ per-shard p99Any endpoint using the queryProjection keyed for the query; partial results; hedgingOwning service team
Cross-shard saga stucksaga.pending_age_seconds p99 ↑Money/inventory in limboReconciler, compensation, alerts at > 60sOwning 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.

OptionWhat WorksWhat BreaksWho Pays
Each service shards independentlyBest local key per serviceA customer spans different shards everywhere; no consistent blast radius; residency = N migrationsEvery team re-solves routing and rebalancing; SRE can't reason about correlated failures
Shared cell architectureBounded blast radius per cell; one tenant-placement system; residency per cellCross-cell features (global search, analytics) need separate pipelines; cell platform is heavyPlatform org (~5–15 FTE); product teams adopt tenant-scoped design
Shared routing/directory library, per-service storesConsistent tenant placement primitives without full cellsPartial alignment onlyPlatform (~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.

ScaleTopologyInfra / MonthPeopleOn-Call Load
Pre-shard (2 TB, 8K writes/s)1 primary + 2 replicas~$10–20K0.5 FTE DB opsLow; failover drills quarterly
Sharded (40 TB, 150K writes/s)16–32 shards × 3 replicas + directory~$150–350KSharding project: 3 FTE × 3 quarters (~$700K one-time); 2–3 FTE ongoingRebalancing, hot-shard pages ~2–4/month
Cells (1 PB, multi-region, 5M tenants)50–200 cells, each with full stack~$2–6MCell 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#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionReversibilityCost to Reverse
Partition keyOne-wayFull reshard; quarters of migration
Number of fixed logical partitionsOne-way-ishReindex / rehash all data (Elasticsearch shards, Kafka partitions)
ID format (shard-aware vs opaque)One-wayIDs are stored and exposed everywhere
Physical node countTwo-wayMove partitions
Directory vs pure hash routingOne-way-ishAdding a directory later requires backfilling mappings and a routing migration
Moving one tenant to a dedicated shardTwo-wayMove it back
Adopting a cell architectureOne-wayOrg-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):

  1. Tenant-scoped data MUST carry tenant_id as the leading partition key component.
  2. Services MUST route through the shared placement directory; no service-private tenant → shard mappings.
  3. Hash-partitioned stores MUST use a fixed logical partition count ≥ 64× the initial node count; hash % nodes routing is prohibited.
  4. Partition ownership changes MUST be epoch-fenced; shards MUST reject writes for partitions they don't own.
  5. 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#

SignalWhat 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)#

ConceptSenior (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#

Related Technologies: DynamoDB · Cassandra · PostgreSQL · Apache Kafka · Redis

  1. Loading the index…