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.

Advertisement

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.

Control plane (cluster state) and data plane (shards) in one Elasticsearch clusterControl plane: master-eligible nodes, voting configuration, cluster stateElected masterpublishes cluster stateMaster-eligiblevotes, acks stateMaster-eligibleor voting_onlypublishCoordinatingnode.roles: [ ]Data node AP0 R1Data node BP1 R0Data node CR0 R1routeprimary P0 forwards each operation to in-sync replicas R0 on B and Crouting tablePer write: seq no + primary termlocal checkpoint per copy, global checkpoint = min over in-sync copiesThe master decides where shards live; it never sits on the document write or search path.
Master-eligible nodes agree on the cluster state and publish it to every node. Coordinating nodes use the routing table to send requests to shard copies. A primary replicates each write to its in-sync replicas and tracks sequence numbers and checkpoints.

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.

Advertisement

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.

WatermarkDefaultMax headroom defaultEffect
Low85%200 GBNo new shards allocated to the node
High90%150 GBShards are moved away from the node
Flood stage95%100 GBindex.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.

  1. The coordinating node looks up the primary for that shard in its cluster state and forwards the request.
  2. 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.
  3. The primary validates and indexes the operation, assigns it a sequence number and writes it to the translog.
  4. 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.
  5. 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 node

Restrict 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

FailureWhat you seeResponse
Lost voting majorityNo elected master; index creation and allocation stopBring voters back; plan decommissions with voting exclusions
initial_master_nodes left in configSeparate clusters bootstrap after a rebuildRemove it after first formation
Flood-stage watermarkWrites rejected with an index blockFree or add disk; the block lifts below the high watermark
Unassigned replicasYellow health_cluster/allocation/explain; check awareness and capacity
Unassigned primaryRed health, writes to that shard failRestore the node or a snapshot; forcing a stale copy loses data
Huge cluster stateSlow publication, master instabilityFewer shards, strict mappings, dedicated masters
Hot nodeOne node's queues and latency much higherCheck 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

  1. Run GET _cat/nodes?v&h=name,node.role,master and confirm you have three master-eligible nodes in different zones and that no data node is master-eligible by accident.
  2. Check that cluster.initial_master_nodes is absent from every running node's configuration.
  3. Set zone attributes and allocation awareness, and decide whether forced awareness fits your capacity.
  4. Alert on disk usage before the low watermark, on yellow and red health, and on unassigned shards with the allocation explain output attached.
  5. Count shards per node and mapped fields per index, and set limits before they become a cluster-state problem.
  6. Write down and rehearse the rolling-restart and master-decommission procedures, including voting exclusions.
Key takeaway: An Elasticsearch cluster is a voting group of master-eligible nodes that publishes one cluster state, and data nodes that hold primary and replica shards described by that state. Majority voting with an odd voting configuration prevents split brain; bootstrapping happens once. Allocation deciders, awareness and disk watermarks decide where shards may live. Writes go through the primary to all in-sync replicas, with primary terms, sequence numbers and checkpoints keeping copies consistent through failover. Operate it by watching voters, disk, shard counts and allocation explanations.