skip to content

How does Elasticsearch decide which shard a document lands on when you index it without a routing value?

level: middleimportance: must knowfreq 70%

answer

  1. No lookup table is involved anywhere
  2. The engine hashes something about the document
  3. The default routing value is the _id
  4. Murmur3 hash, then modulo shard count
  5. Formula bakes in the primary count

basics

~20 s

Elasticsearch hashes the document's routing value, which defaults to its _id, with Murmur3 and reduces that hash modulo the index's routing-shard count to pick exactly one primary shard. Placement is pure arithmetic, computed on any node, never looked up.

solid answer

~40 s

Every index, get, update and delete request carries a routing value; if you do not supply one, Elasticsearch uses the document `_id`. That value is hashed with Murmur3 and mapped onto a single primary shard by `routing_factor = num_routing_shards / num_primary_shards; shard_num = (hash(_routing) % num_routing_shards) / routing_factor`. With `index.number_of_routing_shards` left at the primary count, this collapses to `hash(_routing) % number_of_primary_shards`. Because it is arithmetic, no node needs a per-document directory: any node can compute the target, and a get-by-id is a one-shard operation rather than a broadcast. The price is that the primary shard count is baked into the formula. Change it and every existing document would hash somewhere else, which is why the primary count is fixed at index creation and capacity changes go through split, shrink or reindex instead.

code

bash · 7 lines
bash
# default routing: the _id is the hash input
PUT /orders/_doc/42
{ "tenant": "tenant-7", "total": 19.90 }

# custom routing: tenant-7 is the hash input, so all its docs share one shard
PUT /orders/_doc/42?routing=tenant-7
{ "tenant": "tenant-7", "total": 19.90 }

go deeper

for a junior

Recall that documents are distributed by hashing, that the default hash input is the document id, and that the number of primary shards is chosen once when the index is created.

for a middle

Be ready to write the hash-and-modulo formula and explain why it makes get-by-id a single-shard operation and why the primary count cannot change afterwards.

for a senior

Explain the contract custom routing imposes on every later get, update and delete, and how you would enforce it with a required _routing mapping before a silent 404 reaches production.

for a principal

Own the consequence: placement is a hash function you choose at index-creation time, so the routing input is a schema decision with the same blast radius as the shard count itself.

## Placement is computed, not recorded Elasticsearch keeps no directory saying "document 42 lives on shard 3". A document's home shard is derived from the document itself, so every node in the cluster can work it out independently and instantly. Every document-level request carries a **routing value**. If the request does not set one, the routing value is the document `_id`. Elasticsearch hashes that string with **Murmur3** and reduces the hash to a shard number: ``` routing_factor = num_routing_shards / num_primary_shards shard_num = (hash(_routing) % num_routing_shards) / routing_factor ``` `num_routing_shards` is the static index setting `index.number_of_routing_shards`, fixed when the index is created. When it equals the primary count — the common case — the expression collapses to the classic form `hash(_routing) % number_of_primary_shards`. The resulting `shard_num` names a **primary shard**. The write is sent to that primary's current node, applied there, then replicated to that shard's replicas. Replicas are copies of a specific primary, not independent placement targets: routing chooses the shard, never the copy. ## Why arithmetic beats a lookup table A directory mapping ids to shards would have to be replicated, kept consistent, and would grow linearly with document count — for billions of documents that is a cluster-state and memory problem no one wants. Hashing gives: - **O(1) placement with zero state.** Any coordinating node computes the target from the id in the request. - **Single-shard get/update/delete.** `GET /orders/_doc/42` touches one shard, not all of them. That is why a get by id is dramatically cheaper than a search that matches one document. - **Determinism.** Re-indexing the same id lands on the same shard, so updates and deletes always find the original. - **Even spread.** Murmur3 over high-cardinality ids distributes documents almost uniformly, which is exactly what you want when you have no domain reason to co-locate. ## The consequence everyone is really asking about Because the primary count is an input to the formula, it cannot change on a live index. If you went from 5 primaries to 6, every id would suddenly hash to a different shard and most documents would become unfindable by id. Elasticsearch therefore fixes `number_of_shards` at creation. Replica count, by contrast, is not in the formula at all and can be changed at any time. The split and shrink operations exist precisely because of this arithmetic: `index.number_of_routing_shards` reserves headroom so that primaries can be multiplied by a factor while each document still resolves to the shard that already holds it, and shrink divides in the same controlled way. Both produce a **new index**; neither edits the shard count in place. ## When you supply your own routing value A custom routing value replaces `_id` as the hash input: ``` PUT /orders/_doc/42?routing=tenant-7 ``` Nothing else about the mechanism changes — the same hash, the same modulo. What changes is which documents collide onto the same shard. All of `tenant-7`'s documents now live on one shard, so a search restricted with `?routing=tenant-7` can be answered by that shard alone instead of fanning out. This creates a contract you must honour forever. Once a document is indexed with a custom routing value, **every** subsequent get, update or delete by id must supply the same value, or Elasticsearch will hash the `_id` instead, target the wrong shard, and report the document as missing — a silent 404, not an error. The defence is to declare routing mandatory in the mapping: ``` "_routing": { "required": true } ``` which makes an index request without a routing value fail loudly. When a custom value was used, it is stored and returned as the `_routing` metadata field on the document; the default (the `_id`) is not stored, because it is already there. ## Distribution quality Hashing is only as even as its input. Auto-generated ids give near-perfect balance. Sequential numeric ids also hash well — Murmur3 does not preserve locality. Skew appears when you deliberately reduce cardinality: routing by tenant, country or customer means shard occupancy follows your entity-size distribution, not the hash. One enormous tenant produces one enormous shard, and a shard is the smallest unit the cluster can move, so no rebalancing can fix it. ## What interviewers listen for Strong answers name the routing value explicitly ("defaults to `_id`"), describe hash-then-modulo without pretending to recite exact bit-level details, and connect the formula to the immutable primary count and to the get-by-id fast path. Weak answers describe a load-balanced or round-robin write path, which would make get-by-id impossible without a broadcast.

  • What does index.number_of_routing_shards change about that formula, and why does it exist?
    It replaces the primary count as the modulus and introduces a routing factor that divides the result. Setting it higher than the primary count reserves arithmetic headroom so the index can later be split into a multiple of its primaries while every document still resolves to the shard that already holds it. It is static — you set it at index creation or not at all.
  • If routing decides the shard, what decides whether a read hits the primary or a replica?
    Routing picks the shard; the coordinating node then picks a copy. For searches it balances across primary and replicas using adaptive replica selection, so two identical searches may be served by different copies. Real-time gets by id are also served from any copy by default. The `preference` parameter lets you pin the choice.
  • Can you set a routing value per action inside a bulk request?
    Yes. Each action metadata line in the bulk body accepts its own `routing` field, so one bulk request can write documents belonging to many tenants, each to its own shard. The bulk API groups the actions by target shard internally and sends one sub-request per shard.

saying these in an interview costs you the question

  • Claims a central table records which shard holds each document
  • Says the coordinating node writes to the least-loaded shard
  • Thinks number_of_shards can be raised on a live index
  • Confuses the replica count with the primary count in the formula
  • Believes routing is derived from the index name or timestamp

context