Lesson 21 / 28

Sharding: Splitting Data Across Machines

Understand hash sharding and why changing the shard count is costly.

Each shard holds a slice; queries fan out

When one machine cannot hold or serve all the vectors, the data is sharded: each record is assigned to a shard, usually by hashing its ID (or a tenant key). A query is sent to every shard (or to the shards that can contain matches), each returns its local top-k, and the coordinator merges them into the global top-k. More shards mean more memory and parallelism, but each query touches more machines, so tail latency is set by the slowest shard. Choose the shard key to match access patterns (by tenant keeps one tenant on few shards and makes tenant queries cheap). Resharding a hash-based layout moves most records, so plan the number of shards ahead (or use engines with consistent hashing or many small logical shards that move as units).

Hash sharding and its cost of change, run

I ran this plain-Python (standard library only) example. Hashing 10,000 IDs gives each of 4 shards about 2,500 records and each of 8 shards about 1,250, an even spread. But going from 4 to 5 shards with hash % N moves 7,921 of the 10,000 records, almost 80%, which is why resharding is expensive.

import hashlib
from collections import Counter

def shard_for(doc_id, n_shards):
    return int(hashlib.md5(doc_id.encode()).hexdigest(), 16) % n_shards

ids = [f"doc-{i}" for i in range(10000)]
for n in (4, 8):
    print(n, "shards ->", sorted(Counter(shard_for(i, n) for i in ids).values()))

moved = sum(shard_for(i, 4) != shard_for(i, 5) for i in ids)
print("going from 4 to 5 shards with hash % N moves", moved, "of", len(ids), "documents")

Output:

4 shards -> [2462, 2468, 2519, 2551]
8 shards -> [1204, 1214, 1247, 1254, 1258, 1258, 1272, 1293]
going from 4 to 5 shards with hash % N moves 7921 of 10000 documents

Pick the shard count with growth in mind

Hash-based resharding moves most data, so leave room to grow before you need to.

Quick check: How is a query answered in a sharded vector database?

  • Only the first shard is searched
  • Each shard returns its local top-k and a coordinator merges them
  • The client searches the disks
  • Shards are merged at write time
Answer

Each shard returns its local top-k and a coordinator merges them — Fan-out and merge gives the global nearest neighbours.