Shards — hash(routing) mod N
A primary shard is a complete Lucene index, and Elasticsearch routes documents with shard = hash(routing) % num_primary_shards — which is exactly why num_primary_shards is fixed for the life of the index.
Everything so far fits on one machine. 5 million books at 2 KB of body each is 10 GB of postings — fine for one node. 500 million is not, and the routing decision has to be made before the document even reaches an inverted index.
Scene 05
Shards — hash(routing) mod N
- Watch
- Try it
- Predict
- Capture
Five books, five primary shards. Each PUT computes hash(_id) % 5 and lands on exactly one shard. Watch each one route — the same _id will always land on the same shard, every time, forever.
Highlighted lines are the ones running in the diagram right now.
def routeShard(doc):key = doc.routing or doc._idreturn hash(key) % num_primary_shards
Where this sits in Build a distributed search engine (Elasticsearch / OpenSearch style)
Scene 05 of 12. Each document lives on shard = hash(routing) mod number_of_primary_shards — and that mod is exactly why you cannot change the shard count after you create the index.
Up next. Five primaries on one node loses everything when that node dies, and one node serves at most one node's worth of read traffic. Replicating each primary buys both — but only one of those needs the writer's permission.
All 12 scenes in Build a distributed search engine (Elasticsearch / OpenSearch style) · Every curriculum