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.