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.

Previously

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

  1. Watch
  2. Try it
  3. Predict
  4. Capture
book 1_id=1book 2_id=2book 3_id=3book 4_id=4book 5_id=5shard = hash(_id) % NN = num_primary_shards = 5PARTITIONS · 5shard 0Lucene indexshard 1Lucene indexshard 2Lucene indexshard 3Lucene indexshard 4Lucene index
What to watch for

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.

Continue unlocks when the animation finishes.
Implementation

Highlighted lines are the ones running in the diagram right now.

Router.routeShard
modular hashing — same key, same shard, forever
def routeShard(doc):
key = doc.routing or doc._id
return 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

Built with Arqly
Every scene in Build a distributed search engine (Elasticsearch / OpenSearch style) builds on the one before it.All 12 Build a distributed search engine (Elasticsearch / OpenSearch style) scenes