Lesson 11 / 25
Sharding and Shard Keys
Split data across machines.
Pick a key that spreads load
Sharding splits data across database nodes by a shard key. Hash-based sharding spreads data evenly but makes range queries scatter across shards; range-based sharding keeps ranges together but can create hot spots, such as all new data landing on the latest shard. A good shard key has high cardinality, spreads writes and matches the main query. Cross-shard joins and transactions become expensive, so design queries around the key.
Hash key versus a date-based key, run
I ran this with Python 3.12.3 using only the standard library; inputs are fixed or seeded, so the output is reproducible. 120,000 simulated users with a December signup spike: hashing user_id gives four nearly equal shards (max/avg 1.01), while sharding by signup quarter puts 2.44 times the average on the last shard.
# Choosing a shard key: hash of user_id vs signup month
import hashlib, random
from collections import Counter
rng = random.Random(1)
users = [(uid, rng.choices(range(1, 13), weights=[1] * 11 + [12])[0]) # Dec signup spike
for uid in range(120_000)]
SHARDS = 4
def by_hash(uid, month):
return int(hashlib.sha1(str(uid).encode()).hexdigest(), 16) % SHARDS
def by_month(uid, month):
return (month - 1) // 3 # one quarter per shard
for name, fn in (("hash(user_id)", by_hash), ("signup quarter", by_month)):
counts = Counter(fn(u, m) for u, m in users)
sizes = [counts[s] for s in range(SHARDS)]
print(f"{name:15} {sizes} max/avg = {max(sizes) / (len(users) / SHARDS):.2f}")
Output:
hash(user_id) [30179, 30115, 29969, 29737] max/avg = 1.01 signup quarter [15626, 15866, 15324, 73184] max/avg = 2.44
Plan for hot keys
A celebrity user or viral item can overload one shard regardless of the key; add caching or split hot keys.
Quick check: What is a risk of sharding by creation date?
- It forbids indexes
- Data cannot be queried
- It makes data perfectly even
- New writes all hit the newest shard, creating a hot spot
Answer
New writes all hit the newest shard, creating a hot spot — Range keys can skew load.