skip to content

Why is number_of_shards fixed at creation in Elasticsearch while number_of_replicas can change live?

level: middleimportance: must knowfreq 75%

answer

  1. Where a document lives is computed, not stored
  2. The count sits inside the placement formula
  3. Replicas are copies, not destinations
  4. One is a final setting, the other dynamic

basics

~20 s

Document routing derives the target shard from a hash of the routing value and the primary shard count, so changing that count would send lookups to the wrong shard. Replicas are full copies of primaries and sit outside that formula, so they can be added or removed at any time.

solid answer

~50 s

Elasticsearch picks a document's primary shard from a hash of its routing value (the `_id` unless you supply `_routing`), reduced against the index's primary shard count. Every indexed document is therefore placed by a formula that has the shard count in it. If that count changed, existing documents would suddenly hash to a different shard, so GET, update and delete by id would look in the wrong place and the index would effectively be corrupt. Elasticsearch marks `index.number_of_shards` as a final setting: a `PUT _settings` attempt is rejected outright. A replica is just another copy of an existing primary's Lucene data; it never changes where a document lives, so `index.number_of_replicas` is a dynamic setting that triggers allocation and peer recovery of new copies. When the primary count is genuinely wrong, the escape hatches are reindex into a new index, `_split`, `_shrink`, or rolling over to new indices with a better count.

code

bash · 7 lines
bash
# accepted: replica count is a dynamic setting
PUT /events/_settings
{ "index": { "number_of_replicas": 2 } }

# rejected: final setting [index.number_of_shards] cannot be updated
PUT /events/_settings
{ "index": { "number_of_shards": 10 } }

go deeper

for a junior

Recall the two settings by name: number_of_shards is chosen when the index is created and cannot be changed; number_of_replicas can be updated on a live index at any time.

for a middle

Be ready to explain the routing arithmetic — the shard is computed from a hash of the routing value against the primary count — and why that makes the count a final setting while replicas stay dynamic.

for a senior

Show you have lived with the consequence: name reindex, _split, _shrink and rollover as the real remedies, and say which one you would pick for a live index under write load.

for a principal

Own the design implication: choose an index-per-period pattern so no single sizing decision has to hold forever, and make the shard count a template-level standard rather than a per-index guess.

## The routing formula is the reason When you index a document, the coordinating node has to decide which primary shard owns it. It does that arithmetically, not by lookup: it hashes the routing value — the document `_id` by default, or an explicit `_routing` value if you supplied one — and reduces the hash against the index's shard configuration. In the simplest form this is `shard = hash(_routing) % number_of_primary_shards`. Since Elasticsearch 6/7 the formula also involves `index.number_of_routing_shards`: the hash is taken modulo the routing-shard count and then divided by a routing factor (routing shards divided by primary shards), which is exactly what makes a later `_split` able to move whole groups of routing shards into new primaries without re-hashing. The key property is that the primary shard count is an *input* to that formula. Nothing on disk records "document X lives on shard 3"; the shard is recomputed from the id every time. Change the count and every previously indexed document is now computed to live somewhere it does not. Reads by id would miss, updates would create duplicates on the newly-computed shard, and deletes would silently no-op. There is no cheap in-place fix, because fixing it means physically moving most documents — which is precisely what a reindex does. ## Final settings versus dynamic settings Elasticsearch classifies index settings as *final* (set at creation, never updatable), *static* (updatable only on a closed index) or *dynamic* (updatable on a live index). `index.number_of_shards` and `index.number_of_routing_shards` are final; a `PUT /my-index/_settings` naming them fails with an error saying the setting cannot be updated. `index.number_of_replicas` is dynamic, and so are the settings people tune day to day, such as `index.refresh_interval`. ## Why replicas are different A replica shard is a byte-level copy of a primary's Lucene segments, kept in sync by the replication of every write. It participates in search — the coordinating node may route a shard-level search request to the primary or to any in-sync replica — but it plays no part in deciding where a document belongs. Increasing `number_of_replicas` therefore only means: the master allocates new unassigned replica shards to nodes, and peer recovery copies segments (plus replays translog) from the primary until each copy is in sync. Cluster health goes yellow while those copies are initializing and back to green when they are in sync. Decreasing the count simply deletes copies, which is instantaneous. Replicas are not free, though. Each one multiplies disk usage, must apply the same indexing work as the primary, and consumes heap and file handles like any other shard. They buy availability (a lost node does not lose data) and search throughput (more copies to fan queries across). ## What to do when the primary count is wrong Because the number is frozen, the practical answers are all "make a new index": - **Reindex** into a freshly created index with the right shard count, then swap an alias. Always correct, but re-analyses and re-indexes every document, so it is the most expensive option. - **`_split`** the index into more primaries. The source must be read-only; the target count must be a multiple of the source count and compatible with `number_of_routing_shards`. - **`_shrink`** the index into fewer primaries. The source must be read-only, green, and a copy of every shard must be on a single node; the target count must divide the source count. - **Rollover** — for append-only time-series or log data, stop trying to size one index for all time. Roll to a new backing index periodically (by age, doc count, or primary shard size) and set the shard count in the index template. Each new index gets today's best answer, and yesterday's mistake ages out with retention. ## Interview framing The answer interviewers want is the causal one — routing arithmetic, not "because Elasticsearch says so". A strong candidate then adds the design consequence: because the number is immutable, you either invest in a quick sizing benchmark up front, or you adopt an index-per-period pattern so the decision is revisited automatically and no single index has to be right forever.

  • If you must change the primary shard count, what are your options and how do you pick?
    Reindex into a correctly sized index and swap an alias — always valid, most expensive. `_split` if you need more primaries and the index can go read-only; `_shrink` if you need fewer and can gather a copy of every shard on one node. For append-only time-series data, the cheapest answer is usually neither: fix the index template and let rollover create correctly sized indices going forward.
  • Does custom routing change any of this?
    No — it changes the input to the hash, not the formula. Supplying `_routing` makes all documents sharing that value land on one shard, which makes single-shard searches possible but can create hot shards if one routing value dominates. The shard count is still baked in, so the primary count stays immutable either way.
  • What actually happens in the cluster when you raise number_of_replicas?
    The master adds unassigned replica shards to the routing table and allocates them to eligible nodes. Each new copy runs peer recovery from its primary: segment files are copied, then any translog operations since the copy started are replayed. Health is yellow while copies initialize and returns to green once all are in sync.

The shard number is computed from the document id like a mailbox derived from a house number modulo the number of streets. Add a street and everyone's mail goes to the wrong box, whereas making a photocopy of every mailbox changes nothing about where mail is delivered.

saying these in an interview costs you the question

  • Claims you can update number_of_shards with PUT _settings
  • Says a shard count change only rebalances existing data
  • Confuses replica count with primary count when sizing
  • Thinks Elasticsearch stores a document-to-shard lookup table
  • Believes reindex is unnecessary because Lucene re-hashes automatically

context