Database Sharding Cheat Sheet
Covers horizontal partitioning strategies, shard key selection, routing, and rebalancing for scaling databases across multiple nodes.
Sharding Strategies
Common approaches to splitting data across shards.
- Range-based sharding- Partitions data by key ranges (e.g., user_id 1-1M on shard A); simple but prone to hotspots on sequential keys
- Hash-based sharding- Applies a hash function to the shard key to evenly distribute rows; loses range-query locality
- Directory-based sharding- A lookup service maps each key to its shard, allowing flexible rebalancing at the cost of an extra hop
- Geo-sharding- Partitions by region/location to reduce latency and satisfy data-residency requirements
- Consistent hashing- Maps shards and keys onto a hash ring so adding/removing a shard only remaps a fraction of keys
- Shard key- The column(s) used to determine which shard a row lives on; picking it wrong causes hotspots or cross-shard joins
Hash-Based Shard Routing
Route a request to its shard using a hash of the key.
import hashlibdef get_shard(user_id: str, num_shards: int) -> int: # Consistent hash of the key mod number of shards digest = hashlib.md5(user_id.encode()).hexdigest() return int(digest, 16) % num_shards# Route a query to the correct shard connectionshard_id = get_shard("user_42", num_shards=8)conn = shard_connections[shard_id]conn.execute("SELECT * FROM orders WHERE user_id = %s", ("user_42",))
Consistent Hashing Ring
Minimize key remapping when shards are added or removed.
import bisectimport hashlibclass HashRing: def __init__(self, nodes, vnodes=100): self.ring = {} self.sorted_keys = [] for node in nodes: for i in range(vnodes): key = self._hash(f"{node}:{i}") self.ring[key] = node bisect.insort(self.sorted_keys, key) def _hash(self, key): return int(hashlib.md5(key.encode()).hexdigest(), 16) def get_node(self, key): h = self._hash(key) idx = bisect.bisect(self.sorted_keys, h) % len(self.sorted_keys) return self.ring[self.sorted_keys[idx]]
Common Pitfalls
Issues that surface once a sharded system is in production.
- Cross-shard joins- Joining rows that live on different shards requires app-level fan-out or a scatter-gather query; avoid by denormalizing
- Hotspotting- A poorly chosen key (e.g., monotonically increasing IDs) concentrates writes on one shard
- Rebalancing cost- Adding shards without consistent hashing forces a full data reshuffle; plan capacity ahead of time
- Distributed transactions- Multi-shard writes need two-phase commit or sagas since native ACID transactions don't span shards
- Global secondary indexes- Queries on non-shard-key columns require a separate index service or scatter-gather across all shards
Zero-Downtime Resharding: Dual-Write + Backfill
The standard playbook for moving from N to M shards without an outage.
def write_during_migration(key, value): old_shard = get_shard(key, num_shards=OLD_N) new_shard = get_shard(key, num_shards=NEW_N) write(old_shard, key, value) write(new_shard, key, value) # dual-write so both topologies stay currentdef backfill(old_shards, new_shards): # Copy historical rows written before dual-write was enabled for row in scan_all(old_shards): target = get_shard(row.key, num_shards=NEW_N) upsert(target, row) # idempotent write, safe to re-run# Cutover sequence:# 1. Enable dual-write.# 2. Backfill historical data into new shards.# 3. Verify row counts / checksums match between old and new.# 4. Flip reads to the new shard map.# 5. Stop writing to old shards, decommission.
Scatter-Gather for Cross-Shard Queries
Fan out a query to every shard and merge results in the application layer.
from concurrent.futures import ThreadPoolExecutordef scatter_gather_query(sql, params, shard_connections): def query_one(conn): with conn.cursor() as cur: cur.execute(sql, params) return cur.fetchall() with ThreadPoolExecutor(max_workers=len(shard_connections)) as pool: results = pool.map(query_one, shard_connections.values()) merged = [row for shard_rows in results for row in shard_rows] # Sorting/aggregation (ORDER BY, LIMIT, GROUP BY) must be redone here # since each shard only sorted/limited its own local rows. return sorted(merged, key=lambda r: r["created_at"], reverse=True)
Choosing a Shard Key: Trade-offs
There is no perfect shard key — every choice sacrifices something.
- tenant_id / customer_id- ideal for B2B SaaS: keeps all of a customer's data co-located for fast queries and simple compliance/data-residency boundaries, but risks large-tenant hotspots
- Composite key (tenant_id + entity_id)- combines a low-cardinality routing prefix with a high-cardinality suffix to spread writes within a large tenant across sub-shards
- Reverse-order / salted keys- prefixing a hash or random salt onto a monotonic ID (common in Bigtable/HBase schemas) prevents a single 'hot' region from absorbing all recent writes
- Derived vs. natural keys- deriving the shard key from a hash of a natural key avoids exposing internal routing, but makes manual shard lookup/debugging harder
- Re-shardability- a key chosen without consistent hashing or virtual nodes locks you into painful full-cluster reshuffles when capacity needs to change
Virtual Shards for Cheap Rebalancing
Decouple logical partitions from physical nodes so rebalancing is just moving assignments.
NUM_VIRTUAL_SHARDS = 4096 # fixed forever, far more than any physical node countdef virtual_shard(key: str) -> int: return hash(key) % NUM_VIRTUAL_SHARDS# Small, static mapping updated as physical capacity changesvshard_to_node = { 0: "node-a", 1: "node-a", 2: "node-b", 3: "node-b", # ... 4096 entries}def get_node(key: str) -> str: return vshard_to_node[virtual_shard(key)]# Adding node-c: only reassign a subset of virtual shard entries (e.g.# 0..1365 -> node-c) instead of rehashing every key in the cluster.
Sharding vs. Related Techniques
Terms that get conflated but solve different problems.
- Sharding (horizontal partitioning)- splits rows across independent database instances/nodes to scale writes and total data volume
- Table partitioning- splits rows across sub-tables within a single database instance (e.g. Postgres declarative partitioning); improves query pruning and maintenance, doesn't add write capacity
- Replication- copies the same full dataset to multiple nodes for availability/read scaling; every replica holds all the data, unlike a shard
- Federation / functional partitioning- splits by feature/domain (e.g. users DB, orders DB) rather than by key range or hash
- Sharded + replicated (typical production setup)- each shard is itself a replica set, so you get both horizontal write scaling and per-shard high availability
Pick a shard key with high cardinality and even access patterns (e.g., a hashed user_id), not a monotonically increasing timestamp or auto-increment ID — those funnel all new writes onto the last shard.