skip to content

Routing & the Search Phases

Documents land on a shard by hashing the id, or a custom routing key if you supply one, and a search fans out to every shard in a query phase before fetching only the winners. Interviewers ask what custom routing buys you and what it costs when one tenant is huge.

part ofElasticsearchoverview, primer and where to startread it →
on this pageshow

questions

6

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

open as a page

What happens during the query phase and the fetch phase of an Elasticsearch search?

level: middleimportance: must knowfreq 65%

basics

~20 s

In the query phase each shard runs the search locally and returns only document ids plus sort values or scores; the coordinating node merges these into one globally sorted list. The fetch phase then retrieves the _source of just the winning documents.

open as a page

When does search_type=dfs_query_then_fetch change an Elasticsearch result ordering, and what does it cost?

level: seniorimportance: should knowfreq 45%

basics

~20 s

By default each shard scores using only its own term statistics, so identical documents can score differently per shard. dfs_query_then_fetch adds a preliminary round trip that gathers term and document frequencies from all participating shards so scoring uses global statistics.

open as a page

Why can two identical Elasticsearch searches rank results differently, and how does the preference parameter help?

level: seniorimportance: should knowfreq 40%

basics

~20 s

Successive searches may be served by different copies of the same shard, and a primary and its replica hold different segment layouts and different numbers of not-yet-purged deleted documents, so their local statistics differ slightly. A constant preference value pins requests to the same copies.

open as a page

What does the Elasticsearch shard request cache store, and why do queries containing now never hit it?

level: seniorimportance: should knowfreq 38%

basics

~20 s

It caches per-shard search results keyed on the whole JSON request body, and by default only for requests with size:0 — aggregations, suggestions and the total hit count. A now value resolves to a new instant each request, producing a new key and a guaranteed miss.

open as a page

How would you decide whether to use custom routing for a multi-tenant Elasticsearch index with thousands of tenants?

level: principalimportance: should knowfreq 32%

basics

~20 s

Weigh the fan-out saved against the skew created. Custom routing turns every tenant search into a single-shard search, which multiplies search throughput, but it concentrates each tenant on one shard, so an outsized tenant becomes a hot shard the cluster cannot rebalance away.

open as a page