skip to content

Sharding and Replication

How an index is split into primary shards and copied into replicas, and what that means for the write path, read scaling, and surviving a lost node. Interviewers ask because primary shard count is fixed at index creation, so getting it wrong is expensive to undo.

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

explore

questions

24

What do green, yellow, and red cluster health mean in an Elasticsearch cluster?

level: juniorimportance: must knowfreq 80%

answer

  1. It is about shard copies, nothing else
  2. Two of the three colours still serve all data
  3. Ask which copy is missing: primary or replica
  4. Cluster colour equals the worst index colour

basics

~20 s

Green: every primary and replica shard is assigned. Yellow: all primaries are assigned but at least one replica is not. Red: at least one primary is unassigned, so part of the data is missing from search results and writes to that shard fail.

solid answer

~40 s

Cluster health is purely a **shard-assignment** report. **Green** means every primary and every replica shard has a node to live on. **Yellow** means all primaries are assigned, so all data is searchable and writable, but at least one replica is unassigned — you have full data, reduced redundancy. **Red** means at least one primary is unassigned: that shard's data cannot be searched and indexing into it fails, while the rest of the index still works. Each index gets a colour and the cluster colour is the worst index colour, so one broken index makes the whole cluster red. A single-node development cluster with default settings is permanently yellow, because a replica is never allowed onto the same node as its primary. `GET _cluster/health?level=indices` shows which index is at fault.

code

bash · 1 line
bash
GET _cluster/health?level=indices

go deeper

for a junior

Memorise the three definitions in terms of primary and replica shards, and be ready to explain why your laptop cluster is always yellow.

for a middle

Explain why a replica is never placed on its primary's node, and show how to go from a colour to the individual unassigned shard using the health and _cat APIs.

for a senior

Show the diagnosis path from a red cluster to a root cause, and discuss alerting: which colours page, which merely warn, and why duration matters more than the colour itself.

for a principal

Frame health colour as one signal in an availability SLO. Argue about what redundancy level the business is actually paying for and how many simultaneous node losses each index tier should survive.

## What the colours actually measure Elasticsearch health colour answers exactly one question: *are all the shard copies this cluster is supposed to have actually allocated to a node?* It says nothing directly about CPU, latency, disk, or JVM pressure. Every index is built from a fixed number of primary shards, and each primary can have replica copies. The cluster's job is to find a node for every one of those copies. The colour reports how far it got. ## Green Green means every primary and every replica of every index is assigned to a node. Full data availability, full redundancy: you can lose a node and still have a copy of everything. Green is not a statement about performance — a green cluster can be at 92% disk usage and rejecting requests from queue saturation. ## Yellow Yellow means all primaries are assigned but at least one replica is not. Nothing is missing: every document can be searched and every write succeeds. What you have lost is redundancy — if the node holding an unreplicated primary dies now, that shard goes unavailable and the cluster turns red. Yellow is therefore a "fix it soon", not a "wake up at 3am" (unless it lingers). The most common yellow is not a fault at all. A one-node cluster that creates an index with one replica stays yellow forever, because the allocator refuses to put a replica on the same node as its primary — two copies on one machine buy no redundancy. The fix on a dev box is `number_of_replicas: 0`, not more debugging. ## Red Red means at least one primary shard is unassigned. Documents that hash to that shard are simply absent: searches over the index return partial results (with `_shards.failed` non-zero) and indexing requests routed to it fail. Red is *not* automatically permanent data loss — the usual cause is that the node holding that primary is down or restarting, and the shard comes back when it returns. It becomes real loss only when no copy of that shard exists anywhere any more. A subtlety worth saying out loud in an interview: red is per-shard, not per-cluster. The other shards of the same index, and every other index, keep serving traffic normally. "The cluster is red" is often mis-stated as "the cluster is down". ## How to look it up `GET _cluster/health` gives the colour plus counters: `active_shards`, `unassigned_shards`, `initializing_shards`, `relocating_shards`, `number_of_nodes`. Add `?level=indices` to get a colour per index, or `?level=shards` for per-shard detail. `GET _cat/indices?v&health=red` lists the offending indices. `GET _cat/shards?v&h=index,shard,prirep,state,node,unassigned.reason` lists the individual unassigned shards with a reason code such as `NODE_LEFT`, `INDEX_CREATED`, `ALLOCATION_FAILED`, or `CLUSTER_RECOVERED`. When the reason is not obvious, `GET _cluster/allocation/explain` names the rule that blocked the allocation. ## Typical causes - A data node left the cluster: replicas of its shards go unassigned (yellow), and any primary it held goes unassigned until a replica is promoted (briefly red if there was no replica). - More replicas requested than nodes available: `number_of_replicas: 2` on a two-node cluster is permanently yellow. - Disk watermarks: when a node crosses the high watermark the allocator stops putting shards there, and a cluster with nowhere left to place a copy stays yellow. - Allocation filtering or awareness rules that no node satisfies — for example an index required to sit on a node attribute that no longer exists. - Repeated allocation failures: after `index.allocation.max_retries` attempts the allocator gives up until you ask it to retry. - A restore or a fresh index during initialization is briefly red/yellow while shards initialize; that is normal and self-resolving. ## What the colour is good for As an alert, green-to-yellow is a redundancy alert and yellow-to-red is an availability alert. Both should be scoped by *how long*: transient yellow during a rolling restart is expected and healthy, while yellow for an hour means the allocator cannot place a copy and needs investigation. Alerting on any non-green colour without a duration window produces noise every time a node restarts.

  • A single-node Elasticsearch cluster is yellow and no node has failed. Why?
    The index was created with at least one replica, and Elasticsearch never allocates a replica onto the same node as its primary — two copies on one machine give no redundancy. With only one node, the replica has nowhere to go and stays unassigned. Set `number_of_replicas: 0` for a single-node development cluster, or add a node.
  • Does a red Elasticsearch cluster reject every search request?
    No. Red means at least one primary shard is unassigned; queries against other indices work normally, and a query against the affected index still returns hits from its healthy shards, with the failure reported in the response's `_shards.failed` count. You can require completeness instead by rejecting partial results, but the default is to return what is available.
  • Should you alert on any non-green Elasticsearch cluster health?
    Not on the colour alone. Rolling restarts, node upgrades and new-index initialization all produce short yellow periods. Alert on yellow sustained beyond a few minutes and on red almost immediately, and include `unassigned_shards` in the alert so the on-call engineer sees the scale of the problem.

Think of it as a spare-tyre check, not an engine check: green means every wheel plus a spare, yellow means all four wheels but a missing spare, red means you are driving on three wheels.

saying these in an interview costs you the question

  • Says yellow means the cluster is down or losing writes
  • Claims red always means permanently lost data
  • Treats green as proof of healthy performance and disk
  • Thinks health colour reflects CPU, latency, or heap
  • Does not know cluster colour is the worst index colour

context

open as a page

Why is a document Elasticsearch just acknowledged as indexed sometimes missing from an immediate search?

level: juniorimportance: must knowfreq 78%

basics

~20 s

Elasticsearch search is near-real-time, not real-time. Indexing puts the document in an in-memory buffer and the translog; it only becomes searchable when a refresh turns that buffer into a new searchable Lucene segment, by default about once a second.

open as a page

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

level: middleimportance: must knowfreq 70%

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.

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

Why does an Elasticsearch cluster with thousands of tiny shards perform worse than one with fewer large shards?

level: middleimportance: must knowfreq 70%

basics

~20 s

Every shard is a separate Lucene index with fixed overhead in heap, file handles, merge work and cluster-state metadata, and every search becomes a task per shard that must be scheduled and merged. Those per-shard costs dominate once shards are small, which is why guidance targets tens of gigabytes per shard.

open as a page

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

level: middleimportance: must knowfreq 75%

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.

open as a page

In Elasticsearch, what is the difference between a refresh and a flush?

level: middleimportance: must knowfreq 68%

basics

~20 s

A refresh opens the in-memory buffer as a new Lucene segment so recent documents become searchable, without fsyncing anything. A flush performs a Lucene commit that fsyncs segments to disk and trims the translog, which is about durability and recovery, not visibility.

open as a page

How do _seq_no and _primary_term give Elasticsearch optimistic concurrency control on document updates?

level: middleimportance: must knowfreq 58%

basics

~20 s

Every write stamps the document with the primary's sequence number and the primary term. A client re-sends both as if_seq_no and if_primary_term; the shard applies the write only if they still match the current document, otherwise it returns a 409 version conflict.

open as a page

How do you use GET _cluster/allocation/explain to diagnose an unassigned Elasticsearch shard?

level: seniorimportance: must knowfreq 62%

basics

~20 s

Call it with the index, shard number and whether it is the primary; the response gives unassigned_info explaining why the shard became unassigned, and per-node decider decisions where each NO or THROTTLE names the exact rule blocking allocation on that node.

open as a page

How do you change the number of replicas on a live Elasticsearch index?

level: juniorimportance: should knowfreq 60%

basics

~20 s

Send a PUT to the index's _settings endpoint with index.number_of_replicas. It is a dynamic setting, so it applies immediately without closing the index; new copies are allocated and recovered while health shows yellow until they are in sync.

open as a page

What do the master, data, ingest, and coordinating-only node roles do in an Elasticsearch cluster?

level: middleimportance: should knowfreq 68%

basics

~20 s

In node.roles, master-eligible nodes can be elected to own the cluster state, data nodes hold shards and execute queries on them, ingest nodes run ingest pipelines, and a node with an empty roles list is coordinating-only: it routes requests and merges results.

open as a page

What does Elasticsearch's _split API require of the source index, and how does _shrink differ?

level: middleimportance: should knowfreq 55%

basics

~20 s

Both create a new index rather than resizing in place, and both need the source made read-only first. _split needs a target primary count that is a multiple of the source count and compatible with number_of_routing_shards; _shrink needs a target count that divides the source count, green health, and a copy of every shard gathered on one node.

open as a page

How does refresh=wait_for on an Elasticsearch index request differ from refresh=true?

level: middleimportance: should knowfreq 52%

basics

~20 s

refresh=wait_for holds the response until the next scheduled refresh makes the change visible, forcing no extra work. refresh=true forces an immediate refresh of the affected shards, creating an extra segment per request and much more merge pressure.

open as a page

How do you restart an Elasticsearch data node for maintenance without triggering a full shard reallocation?

level: seniorimportance: should knowfreq 55%

basics

~20 s

Set cluster.routing.allocation.enable to primaries so replicas of the departing node are not rebuilt elsewhere, stop non-essential indexing and flush, restart the node, then set the setting back to all and wait for green. Delayed allocation covers short absences automatically.

open as a page

How does Elasticsearch's voting configuration decide whether a master can be elected after nodes leave?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Elasticsearch maintains a voting configuration: the set of master-eligible nodes whose votes count. An election needs a strict majority of that set, so the cluster survives losing fewer than half of it. Below that it has no master and rejects cluster-state changes.

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

A daily Elasticsearch log index has 20 primary shards holding 400 MB each — how do you fix the sizing?

level: seniorimportance: should knowfreq 50%

basics

~20 s

You cannot change the primary count of existing indices, so fix it in two directions: shrink and force-merge the already-created ones down to one primary, and change the index template plus the rollover trigger so future indices are created with far fewer shards and roll on primary shard size rather than on the calendar.

open as a page

What happens on the primary and its in-sync replicas when Elasticsearch indexes one document?

level: seniorimportance: should knowfreq 62%

basics

~20 s

The coordinating node routes the request to the primary, which validates it, applies it locally with a sequence number and primary term, then forwards it in parallel to every in-sync replica. Once they respond the primary acknowledges; a failed replica is reported to the master and removed from the in-sync set.

open as a page

What does setting index.translog.durability to async cost you when an Elasticsearch node crashes?

level: seniorimportance: should knowfreq 44%

basics

~20 s

With async durability the translog is fsynced on a timer rather than per request, so a node crash or power loss can lose every acknowledged write since the last sync interval. The gain is far fewer fsyncs and higher indexing throughput.

open as a page

How would you lay out an Elasticsearch cluster across three availability zones so losing one zone keeps it writable?

level: principalimportance: should knowfreq 38%

basics

~20 s

Put one master-eligible node in each zone so a majority survives losing one, spread data nodes evenly, tag every node with a zone attribute, enable allocation awareness on that attribute so shard copies land in different zones, and keep at least one replica per index.

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

How do you choose the primary shard count for a new Elasticsearch index whose future size you don't know?

level: principalimportance: should knowfreq 45%

basics

~20 s

Benchmark one shard loaded with representative data and real queries to find the size where latency degrades, then divide the projected data volume by that size. Where growth is genuinely unknown, design so the number never has to be right forever: roll over to new indices sized on shard size.

open as a page