# Partitioning and Sharding — Scalability, Availability & Reliability

Source: https://www.skillbyai.com/en/scalability/d-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.

```python
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 node
```

## Exam 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.

**Quiz:** 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
- [x] 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.
