Lesson 8 / 25
Partitioning and Sharding
Split data across nodes with range, hash and consistent hashing, and avoid hotspots.
When one machine cannot hold the writes
Replication scales reads, but every write still hits one leader. Partitioning (sharding) splits the data so each node owns a subset and handles its writes. Range partitioning assigns key ranges to shards (A–F, G–M): efficient range scans, but sequential keys such as timestamps send all new writes to one shard. Hash partitioning spreads keys evenly by hashing them, but range queries must touch every shard. Consistent hashing places nodes and keys on a ring (often with many virtual nodes per server), so adding or removing a node moves only a small fraction of keys; it is used by Cassandra, DynamoDB-style systems and distributed caches. Choose a shard key that spreads load and matches common queries; a celebrity or hot key (one huge customer, one viral post) can still overload a shard and may need special handling. Cross-shard joins and transactions are expensive, and resharding live data is hard, so many teams start with more logical shards than physical servers.
Hash partitioning versus a naive modulo
Modulo remaps almost every key when the node count changes; consistent hashing does not.
import hashlib, bisect
def h(key: str) -> int:
return int(hashlib.md5(key.encode()).hexdigest(), 16)
# naive: changing len(nodes) from 4 to 5 moves most keys
def node_mod(key, nodes):
return nodes[h(key) % len(nodes)]
class Ring:
def __init__(self, nodes, vnodes=100):
self.ring = sorted((h(f"{n}#{i}"), n) for n in nodes for i in range(vnodes))
self.points = [p for p, _ in self.ring]
def node_for(self, key):
i = bisect.bisect(self.points, h(key)) % len(self.points)
return self.ring[i][1]
ring = Ring(["db1", "db2", "db3", "db4"])
ring.node_for("customer:42")
# adding db5 moves only about 1/5 of keys to the new nodeExam halls by roll number
Assigning halls by surname range (A–F, G–M) is range partitioning: easy to find a range of students, but if half the class is named Sharma one hall overflows. Assigning by a scrambled roll number spreads students evenly, but you must check every hall to list everyone named Sharma.
Quick check: What is the main advantage of consistent hashing over hash-modulo-N?
- It supports SQL joins across shards
- It eliminates hot keys entirely
- It keeps keys in sorted order
- Adding or removing a node moves only a small fraction of keys
Answer
Adding or removing a node moves only a small fraction of keys — Only keys adjacent to the changed node on the ring move.