An Elasticsearch cluster is a set of nodes that agree on two things: who the elected master is, and a single document called the cluster state that says which indices exist, what their mappings are, and which node holds each copy of each shard. Everything else, indexing, searching and recovering, happens on data nodes using that shared map. Most production incidents come from misunderstanding one of these layers: a cluster that cannot elect a master, a shard that will not allocate, a disk that crosses a watermark, or a replica that is quietly dropped from the in-sync set.
This article explains the architecture from first principles. It covers node roles, how master election and the voting configuration avoid split brain, how the cluster state is published, how shards are placed, how a write is replicated and kept consistent with sequence numbers and checkpoints, and how copies recover after failures. Shard sizing, query fan-out and reindexing are covered in designing search at scale, and segments, refresh and scoring in the search system architecture article.
Two planes: cluster state and shards
Think of the cluster as two planes. The control plane is the elected master and the other master-eligible nodes. They maintain the cluster state: node membership, index metadata and mappings, settings, and the routing table that assigns every primary and replica shard to a node. The data plane is the shards themselves. Each shard is a Lucene index on one data node, and a replica is a full copy of a primary on a different node.
The separation matters operationally. The master is not on the request path for documents: a search or index request goes to any node, which uses its local copy of the cluster state to route it to shard copies. A busy master therefore does not slow searches directly, but an overloaded or unstable master stops index creation, mapping updates and shard reallocation, and that eventually becomes an outage.
Node roles and a sensible topology
Each node declares its roles in node.roles. The documented values are master, data and the tier-specific data_content, data_hot, data_warm, data_cold and data_frozen, plus ingest, ml, remote_cluster_client, transform and voting_only. A node with no node.roles setting gets almost all of them. That is convenient for a laptop and bad for production, where you want failures and load in one role to stay away from the others.
# dedicated master-eligible (three of these, one per zone)
node.roles: [ master ]
# hot data nodes holding recent indices
node.roles: [ data_hot, data_content, ingest ]
node.attr.zone: zone-a
# warm data nodes for older, read-mostly indices
node.roles: [ data_warm ]
# coordinating-only: routes requests and reduces results, holds no data
node.roles: [ ]
# a tie-breaker that votes but can never become master
node.roles: [ master, voting_only ]A common medium-sized layout is three small dedicated master nodes, a data tier sized for the shards, and optionally a few coordinating-only nodes when heavy aggregations would otherwise consume data-node heap. A voting-only node is useful as a cheap third vote when you have two larger master-eligible nodes in two zones and a third zone with little capacity.
Discovery, bootstrapping and the voting configuration
Elasticsearch decides things by majority vote. The set of nodes whose votes count is the voting configuration, normally all master-eligible nodes. Electing a master or committing a new cluster state requires responses from more than half of that set. With three voters, any two can proceed and one can fail. With an even number of master-eligible nodes, Elasticsearch leaves one out of the voting configuration so that the set stays odd: if a network partition splits the nodes into equal halves, one half still has a majority, and only that half can elect a master. This is how the cluster avoids split brain, the same idea discussed in the leader election article.
Nodes find each other through discovery.seed_hosts. A brand-new cluster also needs cluster.initial_master_nodes, listing the master-eligible nodes that form the first voting configuration. It is used only for that first bootstrap. Remove it after the cluster has formed, and never set it on nodes joining an existing cluster. If two groups of nodes are each bootstrapped separately with it, you get two independent clusters with the same name, and merging them is not possible without losing data on one side.
As nodes come and go, the voting configuration adjusts automatically while cluster.auto_shrink_voting_configuration is true, which is the default. The hazard is removing many master-eligible nodes at once. If you shut down half or more of the voters before the configuration has shrunk, the remaining nodes cannot form a majority. To decommission master nodes, first exclude them and wait for the change to commit:
POST /_cluster/voting_config_exclusions?node_names=master-old-1,master-old-2
# ... wait until GET /_cluster/state/metadata shows them excluded, then stop them ...
DELETE /_cluster/voting_config_exclusions
How the cluster state is published
Only the elected master changes the cluster state. Each change, such as creating an index, adding a field or moving a shard, produces a new version. The master sends it to all nodes, waits until a majority of the voting configuration has accepted it, then tells everyone to apply it. This two-phase pattern of publish then commit means a master that has lost its majority cannot commit changes, so a deposed master cannot rewrite the routing table behind a new one's back. Nodes apply versions in order and can receive diffs instead of the whole state.
The cost of the cluster state grows with the number of indices, shards and mapped fields, because every node holds a copy and every change is published. This is why mapping explosion, for example from using arbitrary user keys as field names, and very large shard counts hurt even when the data is small. Elasticsearch limits shards per node with cluster.max_shards_per_node, which defaults to 1,000 for non-frozen data nodes. Treat that as a ceiling, not a target. Fewer, larger shards keep the state small and publication fast.
Shard allocation: deciders, awareness and disk watermarks
When a shard needs a home, after index creation, a node failure or a rebalance, the master asks a chain of allocation deciders whether each node may hold it. Every decider can say yes, no or throttle. Some examples: a primary and its replica may not share a node; filtering settings can pin indices to tiers or attributes; awareness spreads copies across zones; disk deciders block nodes that are too full; and throttling limits concurrent recoveries. The default for cluster.routing.allocation.node_concurrent_recoveries is 2.
Awareness turns zones into a placement rule. Tag nodes with node.attr.zone and set cluster.routing.allocation.awareness.attributes: zone, and Elasticsearch keeps copies of a shard in different zones. Forced awareness, set with cluster.routing.allocation.awareness.force.zone.values, stops the survivors from filling up with every missing replica after a whole zone fails, which could otherwise overload them.
| Watermark | Default | Max headroom default | Effect |
|---|---|---|---|
| Low | 85% | 200 GB | No new shards allocated to the node |
| High | 90% | 150 GB | Shards are moved away from the node |
| Flood stage | 95% | 100 GB | index.blocks.read_only_allow_delete on every index with a shard on the node; lifted automatically once usage falls below the high watermark |
The headroom values stop percentages from wasting terabytes on very large disks. When a shard is stuck unassigned, do not guess; ask:
GET /_cluster/allocation/explain
{ "index": "orders-2026.10", "shard": 0, "primary": false }The response lists each node and the decider that said no, for example a same-shard rule, a disk watermark or an awareness rule.
The write path: primary-backup replication
A document goes to the shard chosen by hashing its routing value, which is the _id by default, over the index's primary shards. This is why the number of primary shards is fixed at creation and changes only through split, shrink or reindex. Hash partitioning trade-offs are covered in the sharding article.
- The coordinating node looks up the primary for that shard in its cluster state and forwards the request.
- Before starting, the primary checks
wait_for_active_shards, default 1, meaning only the primary must be active. This is a pre-flight check, not a durability guarantee. - The primary validates and indexes the operation, assigns it a sequence number and writes it to the translog.
- It forwards the operation to every replica in the in-sync set in parallel. Each replica indexes it and records it in its own translog.
- When all in-sync copies have responded, the primary acknowledges the client. By default the translog is fsynced before acknowledging, which is
index.translog.durability: request.
If a replica fails to apply the operation, the primary asks the master to mark that copy stale and remove it from the in-sync set, then acknowledges. The client is not blocked by a sick replica, and a stale copy can never be promoted later. The in-sync allocation IDs live in the cluster state, so this decision survives master failover.
Sequence numbers, primary terms and checkpoints
Three numbers keep copies consistent. The primary term increases each time a new primary is chosen for a shard, so operations from an old primary can be recognised and rejected. The sequence number orders operations within a shard. Each copy tracks its local checkpoint, the highest sequence number below which it has processed everything. The primary tracks the global checkpoint, the minimum local checkpoint across in-sync copies: every operation at or below it is on every in-sync copy.
Worked example: the primary has applied operations 1 to 10. Replica B has 1 to 10, replica C has 1 to 8 and 10, but missed 9. C's local checkpoint is 8, so the global checkpoint is 8. If the primary's node dies now, the master promotes B. B increments the primary term and resynchronises C by replaying operations above the global checkpoint, which delivers 9 to C. Any operation that the old primary applied but never sent to anyone, operation 11 for example, was never acknowledged to the client and is discarded. Acknowledged writes survive as long as one in-sync copy survives.
The read path in one paragraph
A search goes to a coordinating node, which picks one copy of each shard. Adaptive replica selection prefers copies on nodes with lower recent latency and shorter queues. The query phase returns top document IDs and scores from each shard, the coordinating node merges them, and the fetch phase retrieves the winning documents. Reads only see documents after a refresh, so search is near-real-time, while a GET by ID is real-time by default. Replica routing in general is discussed in the read-replica routing article.
Node loss, recovery and rolling restarts
When a data node leaves, the master promotes replicas for any primaries it held, which is quick. Then it waits index.unassigned.node_left.delayed_timeout, default one minute, before rebuilding missing replicas elsewhere, because a restarting node usually returns and copying whole shards across the network would be wasted work.
Rebuilding a replica is a peer recovery from the primary. If the replica already holds most of the data and the primary still retains the operations it is missing, which it tracks with retention leases, the recovery replays only those operations. Otherwise it copies segment files, then replays operations written during the copy. Recoveries are throttled per node, so a large node loss can take hours to heal. During that time the cluster health is yellow and a second failure can make it red.
A rolling restart follows from these mechanics:
PUT _cluster/settings
{ "persistent": { "cluster.routing.allocation.enable": "primaries" } }
POST _flush
# restart one node, wait for it to rejoin (GET _cat/nodes)
PUT _cluster/settings
{ "persistent": { "cluster.routing.allocation.enable": null } }
# wait for green (GET _cluster/health?wait_for_status=green), then the next nodeRestrict allocation so replicas are not rebuilt while the node is away, restart it, re-enable allocation and wait for green before the next node. Restart master-eligible nodes one at a time so a majority always remains.
Failure modes and trade-offs
| Failure | What you see | Response |
|---|---|---|
| Lost voting majority | No elected master; index creation and allocation stop | Bring voters back; plan decommissions with voting exclusions |
initial_master_nodes left in config | Separate clusters bootstrap after a rebuild | Remove it after first formation |
| Flood-stage watermark | Writes rejected with an index block | Free or add disk; the block lifts below the high watermark |
| Unassigned replicas | Yellow health | _cluster/allocation/explain; check awareness and capacity |
| Unassigned primary | Red health, writes to that shard fail | Restore the node or a snapshot; forcing a stale copy loses data |
| Huge cluster state | Slow publication, master instability | Fewer shards, strict mappings, dedicated masters |
| Hot node | One node's queues and latency much higher | Check routing skew and shard balance |
Two trade-offs deserve explicit choices. More replicas buy read throughput and failure tolerance but multiply indexing work and disk, since every copy indexes every document. Synchronous replication to in-sync copies keeps acknowledged writes safe but means write latency is set by the slowest healthy copy. Spreading a cluster across distant regions makes that latency painful, which is why cross-region setups usually use separate clusters with asynchronous replication.
What to do next
- Run
GET _cat/nodes?v&h=name,node.role,masterand confirm you have three master-eligible nodes in different zones and that no data node is master-eligible by accident. - Check that
cluster.initial_master_nodesis absent from every running node's configuration. - Set zone attributes and allocation awareness, and decide whether forced awareness fits your capacity.
- Alert on disk usage before the low watermark, on yellow and red health, and on unassigned shards with the allocation explain output attached.
- Count shards per node and mapped fields per index, and set limits before they become a cluster-state problem.
- Write down and rehearse the rolling-restart and master-decommission procedures, including voting exclusions.