Core Concept

Database Sharding and Partitioning

Sharding splits data across nodes to scale writes and storage, but your shard key determines which queries stay fast and which become expensive scatter-gather scans.


1. What It Is

You reach for sharding when a single primary can't keep up with write QPS or disk growth β€” but it's expensive to undo, so we only shard after caches and read replicas are exhausted.

What:

Horizontal partitioning (sharding) splits single tables horizontally into separate database servers; vertical partitioning splits table columns into separate tables.

Primary purpose:

Scale write throughput and storage capacity beyond the physical hardware limits of a single master server.

Usually used for:

High-volume transactional databases, heavy social timeline layers, and multi-tenant SaaS backends.

2. Core Mental Model

The shard key is the most important decision you'll make β€” it determines which queries hit one node and which fan out to every shard:

πŸ”‘ The Sharding Key Rule

Select a key (e.g. `user_id` or `tenant_id`) that routes over 90% of your common queries to a single shard, entirely avoiding cross-shard lookups.

βš–οΈ Even Distribution

Choose keys with high cardinality to distribute storage and write QPS uniformly across your database node pool.

🧬 Application Routing

Move routing logic either into an application middleware library (e.g., Vitess) or deploy a dedicated proxy tier (e.g., Citus).

In the room

Interviewers love asking "how would you query across shards?" If your shard key is user_id but the product needs a global leaderboard, you have a problem. Name the shard key early and explain which queries stay single-shard vs which need a secondary index or denormalized table.

3. Why It Matters in HLD

Sharding is the last resort after vertical scaling, caching, and read replicas β€” but when we need it, the shard key is everything. Three lenses:

Needed When:

Database write volume outgrows physical disk write I/O speeds, or total storage requirements exceed single SSD arrays.

Avoids:

Primary database CPU locks during writes, slow sequential table scans, index cache churns, and massive hardware migration costs.

Optimizes For:

Write concurrency scaling, distributed storage limits, data locality latencies, and transaction boundary isolations.

4. Architecture & Data Flow

Walk partition strategies as interview steps. Step 1 β€” Range: contiguous key ranges per shard β€” great for time-series, risky for hot ranges. Step 2 β€” Hash: hash(key) mod N or consistent hash ring β€” even spread, cross-range queries expensive. Step 3 β€” Directory: lookup service maps keys to shards β€” flexible moves, extra hop. Step 4 β€” Geo: pin data to regions for GDPR or latency. Step 5 β€” Name the hot query: pick the strategy that keeps 90%+ of queries single-shard.

Loading...

In the room

Name your shard key in the first five minutes. If it is user_id but the product needs a global leaderboard, say how you would denormalize or scatter-gather β€” do not wait for the interviewer to catch the mismatch.

5. Key Characteristics

Each strategy determines whether queries hit one shard or fan out to all of them β€” we compare before committing:

  • Partition strategy determines whether queries hit one shard or fan out to all of them:
StrategyRouting ModelPrimary BenefitSevere Cost
Range-BasedMap key intervals (e.g. A-E, F-L) to specific shard nodes.
  • Trivial to implement
  • ideal for clustered range queries.
Creates severe hotspots (e.g. all updates on active suffix keys).
Hash-BasedApply hash function to partition key modulo node count: hash(key) % N.Ensures uniform data distribution across nodes.Adding/removing nodes remaps almost 100% of keys (requires Consistent Hashing).
Directory (Lookup)Maintain a routing table mapping keys or ranges to shard IDs.Flexible rebalancing and hot-key relocation without changing hash function.
  • Directory service becomes a coordination bottleneck
  • must stay highly available.
Geo-LocatedRoute user tables to nodes inside their geographic region.
  • Sub-millisecond latency
  • solves local compliance laws.
Complex cross-region backups and failovers.

6. Strategic Tradeoffs

Horizontal partition buys write scale at the cost of cross-shard complexity β€” we state both:

BenefitCost
Horizontal Scalability (bypasses the physical scale limit of a single server, scaling database QPS linearly)
  • No Distributed Joins (joining tables across distinct shards is extremely slow
  • forces application-layer joins)
Small Blast Radius (if shard node A crashes, only 10% of users are degraded, keeping the remaining pool alive)Loss of ACID Atomicity (transactions span multiple shards, requiring complex 2-Phase Commit protocols)

7. Failure / Bottleneck Awareness

Hot shards and scatter-gather are the classic sharding interview traps β€” we volunteer fixes:

πŸ”₯ Celebrity Hot Shard

Problem: Sharding by `user_id` works uniformly until a celebrity account (e.g. 100M followers) executes massive read/write activities. The destination shard node hosting that celebrity collapses under the QPS load, while other shards remain idle.

Mitigation: Implement a hybrid key model: append random salt suffixes to hot user keys (e.g. `celebrity_123:salt_0` to `celebrity_123:salt_9`) to spread their data across multiple shards.

🐒 The Scatter-Gather Query Penalty

Problem: Executing queries without specifying the partition sharding key forces the database coordinator proxy to query all shards concurrently and merge the results. This destroys throughput and maximizes tail latency.

Mitigation: Ensure all high-volume transactional queries include the sharding key, or maintain secondary indexes mapped in a dedicated lookup directory table.

8. Common HLD Usage

These production shard keys show how real teams co-locate data with access patterns:

Production SystemSharding Key SelectionArchitectural Rationale
Twitter Timeline Shardsuser_id (Hash-sharded)
  • Queries pull timeline segments per user
  • sharding by user_id keeps reads bound to a single shard.
E-Commerce Ledgertenant_id (Lookup-sharded)Multi-tenant B2B platforms isolate customers completely, routing all tenant requests to dedicated databases.
URL Shortenershort_code (Directory lookup)A directory table maps vanity codes to shard IDs so hot codes can be relocated without re-hashing the key space.

9. Decision Signals

Shard only when measured write QPS or disk size exceeds single-node limits β€” not because DAU sounds big:

🎯 Think Sharding When:
  • Your database storage volume exceeds 2 Terabytes, or write requirements exceed 10,000 write QPS.
  • You can easily isolate users/tenants data into disjoint sets with zero overlap requirements.
  • You need to geographical locate data in specific datacenters to comply with data privacy policies (GDPR).

11. Deep Dive (Optional)

Vertical Partitioning vs Horizontal Sharding

It is critical to distinguish between these two sizing dimensions:

1. Horizontal Sharding

Splits table **rows** across multiple database servers. The schema remains identical across all shards, but each shard stores a unique subset of rows. Ideal for scaling QPS and disk capacities.

2. Vertical Partitioning

Splits table **columns** into separate tables stored on different disks or servers. For example, moving a large text column (`user_bio`) or blob column (`profile_photo_blob`) out of the core `users` table into a secondary `user_profiles` table.

**Benefit**: Reduces the size of the primary table's memory-mapped pages, significantly boosting scan speeds for common narrow queries (e.g. `SELECT username FROM users`).

When Sharding Is Premature

Sharding is an operational tax β€” cross-shard JOINs disappear, distributed transactions get harder, and rebalancing shards is a multi-week project. Before partitioning horizontally, exhaust vertical scaling and simpler wins: read replicas, connection pooling, query/index tuning, caching hot keys, and archiving cold data. A single well-tuned PostgreSQL instance often carries millions of DAU; interviewers reward candidates who shard only when measured QPS or disk size forces it, not because the problem mentions "scale."

Cross-Shard Operations

Once sharded, queries that span keys on different partitions become expensive:

  • Scatter-gather: Fan out to all shards, merge results in the app layer β€” acceptable for rare admin reports, deadly for user-facing search at high QPS.
  • Global secondary indexes: Maintain a separate index shard keyed differently (e.g., email β†’ user_id lookup) with async replication β€” adds write amplification.
  • Co-locate related data: Choose shard key so hot queries hit one partition β€” user_id shards user posts, orders, and settings together.

Design the shard key around your hottest access path, not around entity nouns in isolation.

πŸ’¬Review

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...