How Sharding Scales Relational Databases
How horizontal partitioning distributes database rows across multiple physical nodes using hash, range, and directory-based sharding strategies.
Last reviewed: July 27, 2026
Sharding splits a database's rows across multiple separate database instances based on a partition key (like customer ID or region), so each shard holds a subset of the data and handles a subset of the load. It lets a database scale write throughput and storage past what a single machine can handle, at the cost of losing easy cross-shard joins and transactions.
Introduction
A single PostgreSQL server can handle millions of rows and thousands of queries per second — until it can’t. When write throughput exceeds what a single disk can handle, or when the working set exceeds available RAM, a single node becomes a bottleneck regardless of vertical scaling. Database sharding is the practice of horizontally partitioning data across multiple independent database instances (shards), each responsible for a subset of the data. When a query arrives, a routing layer determines which shard(s) to query.
Step-by-Step: Sharding Architectures
flowchart TD
A["1. Choosing a Shard Key"]
B["2. Hash Sharding"]
C["3. Range Sharding"]
D["4. Application-Level Routing"]
E["5. Cross-Shard Operations and Their Costs"]
A --> B
B --> C
C --> D
D --> E
Step 1: Choosing a Shard Key
The shard key determines how rows are distributed across shards. The ideal shard key:
- Has high cardinality (many unique values)
- Produces even distribution (no hot shards)
- Appears in most queries (avoid cross-shard joins)
- Is immutable (changing a shard key requires re-sharding)
-- Common shard key choices:
user_id -- Good: high cardinality, appears in most user queries
tenant_id -- Good: multi-tenant SaaS, co-locates tenant data
created_at -- Bad: time-series causes hot shard (all writes go to latest)
email -- Bad: requires normalization, lookup overhead
Step 2: Hash Sharding
Apply a hash function to the shard key modulo the number of shards:
def get_shard(user_id: int, num_shards: int = 4) -> int:
return hash(user_id) % num_shards
# user_id=1001 → shard 1
# user_id=1002 → shard 3
# user_id=1003 → shard 0
Advantage: Even distribution regardless of key patterns. Disadvantage: Re-sharding (changing num_shards) requires moving ~75% of data. Mitigated by consistent hashing (only moves 1/N of data when adding a shard).
Step 3: Range Sharding
Assign ranges of the shard key to specific shards:
Shard 0: user_id 1 – 1,000,000
Shard 1: user_id 1,000,001 – 2,000,000
Shard 2: user_id 2,000,001 – 3,000,000
Advantage: Range queries are shard-local (e.g., “all users registered this month”). Disadvantage: New user signups always hit the last shard (write hot spot) unless shards are pre-split.
Step 4: Application-Level Routing
import hashlib
SHARD_MAP = {
0: "postgresql://shard0.db.internal:5432/app",
1: "postgresql://shard1.db.internal:5432/app",
2: "postgresql://shard2.db.internal:5432/app",
3: "postgresql://shard3.db.internal:5432/app",
}
def get_db_connection(user_id: int):
shard_id = int(hashlib.md5(str(user_id).encode()).hexdigest(), 16) % len(SHARD_MAP)
return connect(SHARD_MAP[shard_id])
# Query routing
def get_user(user_id: int):
conn = get_db_connection(user_id)
return conn.execute("SELECT * FROM users WHERE id = %s", [user_id])
# Cross-shard aggregation (fan-out)
def get_total_users():
total = 0
for shard_conn in [connect(url) for url in SHARD_MAP.values()]:
total += shard_conn.execute("SELECT COUNT(*) FROM users").fetchone()[0]
return total
Step 5: Cross-Shard Operations and Their Costs
| Operation | Same-Shard | Cross-Shard |
|---|---|---|
| Single row lookup | O(1) fast | Route to correct shard |
| Range query (shard key) | Single shard | Fan-out to all shards, merge |
| JOIN across shard keys | Shard-local | Requires application-level join (very expensive) |
| Transactions | ACID guaranteed | Distributed transaction required (2PC, expensive) |
| Aggregations (COUNT, SUM) | Single shard | Fan-out query, reduce results |
Cross-shard JOINs are the primary reason to design schema around co-locating related data on the same shard (e.g., store all of a user’s orders on the user’s shard).
Key Takeaways
- Shard key selection is the most critical architectural decision — a bad shard key causes hot spots, expensive cross-shard queries, or re-sharding pain.
- Hash sharding provides even distribution; range sharding enables efficient range scans at the cost of potential write hot spots.
- Consistent hashing reduces data movement when adding shards from ~75% (modulo) to ~1/N.
- Cross-shard JOINs and transactions are expensive — design schema to co-locate related data on the same shard.
- Managed solutions (CockroachDB, PlanetScale, Vitess, Amazon Aurora Sharding) handle routing, re-sharding, and distributed transactions — prefer them over hand-rolled sharding.
Common questions
How is sharding different from a read replica?
A read replica holds a full copy of the same data and only offloads reads. Sharding splits the data itself across instances — each shard holds a different slice of rows, and both reads and writes for a given row go to its specific shard.
What makes choosing a shard key hard?
A bad shard key creates "hot" shards — if you shard by customer and one customer is far larger than the rest, that shard takes disproportionate load. A good key distributes both data volume and query patterns evenly, which usually means understanding your actual access patterns before picking one.
Why do cross-shard queries become expensive?
A query that needs data from multiple shards (a join across two sharded tables, for example) can't be satisfied by a single database instance anymore — the application or a query router has to fan the query out to multiple shards and merge results itself.
Historical figures, architectures, and capabilities are for informational purposes only. Not technical, professional, legal, or financial advice. Sources: Research papers, developer documentation.