# Sharding: Splitting Data Across Machines — Vector Databases

Source: https://www.geekswithgeeks.com/en/vector-databases/o-shard

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

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

**Quiz:** How is a query answered in a sharded vector database?

- [ ] Only the first shard is searched
- [x] 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.
