Sharding, Replicas, Routing & Cluster Health
An index is split into primary shards with zero or more replicas. Writes route by a key, searches scatter and gather, and cluster health is green, yellow, or red based on whether those shards are allocated.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
What is a primary shard responsible for?
Answer
It is the write authority for its slice of the index. Searches can also run on it. The coordinating node is not a data owner.
L2
What do replicas buy, and what do they cost?
Answer
They buy another searchable copy and a failover candidate. Every index write is replicated, so they add disk, CPU, and write amplification.
L3
Why is a single-node cluster yellow when replicas are 1?
Answer
The primary is allocated. There is nowhere legal to put the replica, because a replica should not sit on the same node as its primary. Dev clusters set replicas to 0 or add a node.
L4
Can you change the primary shard count in place?
Answer
Not as a setting update. number_of_shards is fixed at creation. Later you reindex, or you split or shrink under those APIs' constraints. Replica count can change on a live index.
L5
When is custom routing a good idea?
Answer
When a query should hit one shard, such as all documents for a tenant you can afford to isolate. The same routing value must be sent on get, update, and delete or you miss the document.
L6
What does a disk watermark do?
Answer
Low, high, and flood-stage watermarks stop the cluster from assigning shards, or from writing, when a disk fills. A full disk can leave shards unassigned and turn health yellow or red.
L7
How do tiny shards hurt relevance?
Answer
BM25 uses per-shard statistics. A tiny shard has noisy IDF. Oversharding is both an ops problem and a scoring problem.
Failure modes
Oversharding
Thousands of tiny shards inflate cluster state, heap, and recovery work.
Hot routing key
One user or tenant owns a shard that runs hot while other nodes idle.
Replica on a one-node cluster
Health stays yellow because the replica cannot allocate.
Mismatched routing on read
A get with the wrong routing value looks at the wrong shard and returns not found.
Flood-stage watermark
The node stops writes. Indices on that node become read-only until you free disk.
Forced awareness with no room
Zone rules leave replicas unassigned forever, so the cluster never returns to green.
Misconceptions
Replicas make indexing faster.
They copy the write. Indexing throughput usually drops as replica count rises.
You can raise number_of_shards on a live index the way you raise replicas.
Primary count is a creation-time choice. Replica count is the dynamic one.
The coordinating node holds the data.
It parses, fans out, and merges. Data nodes hold the shards. Any node can coordinate a request.
Interviewer traps
Reciting consistent-hashing ring math for Lucene routing.
Say the default is a hash of _id onto the primary shards of that index. Stay on shard size, replicas, and health.
Treating yellow as healthy because searches return.
Yellow means you are one failure away from losing a primary's redundancy. Say what is unassigned.
How many primary shards for a 100 GB catalog?
Prefer
A handful of shards near the size band
About 100 GB at a 30 GB target is a few primaries, not one shard per day and not one shard per node you hope to buy later. Replicas start at 1 once you have a second data node.
- Cluster state stays small.
- Recovery moves a bounded amount of data.
- BM25 statistics are less noisy than on tiny shards.
- You still leave a split or reindex path for later growth.
Alternative
One primary per future node, plus a replica, on day one
Empty capacity becomes thousands of shards as indices multiply. Searches pay overhead on every shard they touch.
- Heap goes to shard metadata.
- A single-node trial stays yellow if replicas are 1.
- Tiny shards skew IDF.
- Changing your mind is a reindex, so the extra shards are sticky.
Where one index request goes
The coordinating node is a router. The primary is the authority.
- 1
Any node can coordinate
It parses the request. It does not have to own the shard. - 2
Hash the routing key
The default key is the document id. A custom routing value replaces it. The hash selects a primary. - 3
Primary writes, replicas copy
The write is acknowledged after the copies you required say so. More replicas mean more copies. - 4
Search scatters
The coordinator asks a copy of each relevant shard, then merges hits. A custom route can skip shards. - 5
Allocation fails closed
No home for a primary turns the index red. No home for a replica turns it yellow.
Overview
An index is split into primary shards. Each primary has zero or more replicas. A coordinating node routes writes and scatters searches. Cluster health is the allocation report: green, yellow, or red.
Wrong shard counts are a top production incident. So are routing hotspots and disks that cross a watermark. This page is that ops story. Analyzer details stay on the mapping lesson. BM25's formula stays on the query lesson, except for one fact: scores use per-shard statistics, so shard size is also a relevance choice.
Primaries, replicas, coordinators
| Role | Job | What you gain | What you pay |
|---|---|---|---|
| Primary shard | Writes for that slice, and searches | A horizontal cut of the index | Too many means overhead. Too few means huge shards and hot spots |
| Replica shard | A copy of one primary | Read scale and failover | Disk, CPU, and write amplification |
| Coordinating node | Parse, fan out, merge | You can scale the gateway path | A bad aggregation still taxes the data nodes |
Flow
- 1
Coordinating node
- nextPrimary shard 0
- nextPrimary shard 1
- 2
Primary shard 0
- nextReplica of shard 0
- 3
Primary shard 1
- nextReplica of shard 1
- 4
Replica of shard 0
- 5
Replica of shard 1
Lesson map
Sharding, Replicas, Routing & Cluster Health
An index is split into primary shards with zero or more replicas. Writes route by a key, searches scatter and gather, and cluster health is green, yellow, or red based on whether those shards are allocated.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB master["Master"] n1["Data node A"] n2["Data node B"] n1 -->|1. Node B| master master -->|2. Promote a| n1
The replica of shard 0 should live on a different node, and ideally a different zone, than primary 0. Allocation awareness and forced awareness exist for that. A forced rule with no eligible node leaves the replica unassigned. The cluster can sit yellow until you fix the rule or add capacity.
Sizing that survives an interview
Elastic's own guidance is a band, not a law: shards of roughly 10 GB to 50 GB, and on the order of 200 million documents per shard, work for many workloads. Oversharding is the failure mode they name. A thousand 5 GB shards cost more than a smaller set of larger ones.
- Prefer fewer larger shards over a shard per tenant or per day when those shards stay tiny.
- Leave RAM for the operating system page cache. Lucene reads files. A heap that consumes the machine steals that cache. Teams often keep heap in the low tens of GB per node and measure.
index.number_of_shardsis set at index creation. Replica count can change later.- To change primary count you reindex, or you use the split or shrink APIs. Shrink requires the target count to be a factor of the source count, and a copy of every shard must sit on one node first.
- There is also a cluster shard limit,
cluster.max_shards_per_node, defaulting to 1000 non-frozen shards per data node counting primaries and replicas. It exists so a bug cannot create an unbounded shard storm.
Time-based indices multiply this math. Five primaries times two copies times 400 daily indices is a shard explosion. Rollover and retention, on the next page, are how you stop creating a new index that never gets large enough.
Lucene's on-disk unit inside a shard is a segment. That layout is not an OLTP B-tree. Crash-recovery families for Postgres and RocksDB are Database storage engines — WAL, B-trees, and LSM trees. Here you only need: a shard is one Lucene index, and its segment count is an indexing concern.
Routing and hotspots
| Choice | Behavior | Risk |
|---|---|---|
Default _id | Hash across primaries | Usually even, and a get by id just works |
| Custom routing, such as a tenant id | Related documents share a shard | A whale fills one shard |
| One index per day, many shards | Isolation by time | Tiny shards and a huge cluster state if ILM never rolls them up |
Hotspot symptom: two nodes at high CPU or disk while the others idle. Check shard sizes, the routing values, and whether a heavy aggregation is hitting every shard anyway. A search with no routing key still scatters.
Custom routing is part of the document's identity. Index, get, update, and delete must send the same value. A get that hashes a different key looks at the wrong shard and reports the document missing.
Health colors
| Color | Meaning | Typical cause |
|---|---|---|
| Green | Every primary and every replica is allocated | The assignment you asked for exists |
| Yellow | Primaries are allocated, at least one replica is not | Single-node cluster with replicas set to 1, or a disk or awareness rule |
| Red | At least one primary is missing | Node loss before a replica can be promoted, or an allocation failure for that primary |
Sequence
- 1
Data node A → Master
1. Node B stopped reporting
- 2
Master → Data node A
2. Promote a replica if this node has one
- 3
Master
3. Yellow if a replica is missing
- 4
Master
4. Red if a primary has no copy left
Yellow on a laptop with number_of_replicas at 1 is expected. Searches still run on the primary. You have no failover copy. Production treats yellow as a page, not as a mood.
Watermarks
Disk watermarks are the other allocation story:
- The low watermark stops new shards from landing on a node that is too full.
- The high watermark tries to move shards off that node.
- The flood stage refuses writes on indices that have a shard on that node, often by marking them read-only.
A cluster that "randomly" stops indexing is often a disk watermark, not a mapping mystery. Free space, then clear the read-only block, then ask why the disk filled (unbounded indices, replica copies, or snapshots writing locally).
A planning helper
The helper is a ceiling division into a target shard size, capped so a whiteboard answer cannot recommend 200 primaries. It is not a capacity model.
Press Run. Snippets must be self-contained — no network, files, or native modules.
100 GB divided by 30 GB is 3.33, which ceilings to 4. 250 GB ceilings to 9, still under the cap. Say the band out loud, then look at document count and query cost before you lock the number.
Press Run. Snippets must be self-contained — no network, files, or native modules.
Interview Q&A
Do replicas speed up indexing?
Answer
No. Each write is copied to the replicas. Replicas help search throughput and availability. They add write cost and disk.
Why is a single-node cluster yellow?
Answer
Replicas are configured, and a replica will not allocate onto the node that already holds its primary. Set replicas to 0 for a solo dev node, or add a data node.
Can you change the primary shard count in place?
Answer
Not with an update setting. Plan it at index creation. Afterwards use reindex, split, or shrink. Shrink needs a factor of the current count. Replica count is the knob you can turn on a live index.
What is a coordinating node?
Answer
The node that accepted the request. It routes the write to the primary and scatters the search, then merges. It is a role for that request, not a special data-free tier you must buy, though dedicated coordinating nodes exist.
When does custom routing fail a get?
Answer
When the get omits the routing value or sends a different one. The hash selects a different shard, and the document is not there. Treat the routing value as part of the id.
What do you check when two nodes are hot and the rest are idle?
Answer
Shard sizes and which nodes own them, the routing key distribution, and whether the hot query is an aggregation that still touches every shard. Then look at disk watermarks so a full node is not also refusing work.
How do shards change relevance?
Answer
BM25's IDF is per shard by default. Evenly sized shards keep those statistics closer. A global-frequency search type exists and costs a round trip. The query lesson owns the formula.
Green versus yellow versus red?
Answer
Green: every copy you asked for is allocated. Yellow: every primary is allocated and at least one replica is not. Red: at least one primary is missing, so a slice of the index cannot be served.
Pitfalls
Three data nodes, five primaries, one replica. Remove one node. Say whether you expect yellow or red, which copies can be promoted, and what a flood-stage watermark would do if the remaining disks are nearly full.