What happens on the primary and its in-sync replicas when Elasticsearch indexes one document?
answer
- one copy assigns the ordering, others follow
- the set of copies is not just the replica count
- a sick replica must not block writes
- but it must be marked before acknowledging
- checkpoints decide how cheap recovery is
basics
~20 sThe 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.
solid answer
~50 sAny node can receive the write and acts as coordinating node. It hashes the routing value to pick the shard and forwards the request to that shard's primary. The primary validates the request, checks `wait_for_active_shards` before proceeding, applies the operation locally — assigning the next `_seq_no` and stamping the current `_primary_term`, writing to the buffer and the translog — and then forwards it concurrently to every copy in the **in-sync allocation IDs** set held in the cluster state. Each replica applies the same operation with the same sequence number and term. When all in-sync copies have responded, the primary reports success back to the coordinating node. If a replica fails or times out, the primary asks the master to remove that copy from the in-sync set and the write still succeeds — Elasticsearch never lets a sick replica block writes, but it does ensure the copy is marked stale so it can never later be promoted as if it were current.
code
bash · 9 lines# Require at least 2 active copies of the shard before the write proceeds
PUT /orders/_doc/abc-123?wait_for_active_shards=2
{
"customer": "c-99",
"total": 42.50
}
# Inspect the in-sync copies and their recovery state
GET /_cat/shards/orders?v&h=index,shard,prirep,state,nodego deeper
Recall the shape: coordinating node routes to the primary, the primary applies the write, then replicas apply the same operation before the client hears success.
Explain the sequence number and primary term assignment, the parallel replication to in-sync copies, and the fact that durability comes from the translog on each copy.
Be able to reason about failures: what happens when a replica or the primary dies mid-write, why the master is told before the acknowledgement, and why client writes should be idempotent.
Own the guarantee you publish to application teams — what an acknowledgement means, the ambiguity window on timeouts, and the idempotency conventions services must follow when writing to the cluster.
## The path of one write 1. **Coordinating node.** The client's request lands on some node in the cluster. That node hashes the routing value (the document `_id` by default, or a custom `routing` value) to determine which shard owns the document, then looks up the current location of that shard's primary in the cluster state and forwards the request there. 2. **Pre-flight check on the primary.** Before doing any work the primary consults `wait_for_active_shards`, which defaults to `1` — meaning the primary itself. If the required number of active copies is not available, the operation waits (up to a timeout) rather than proceeding. This is a *pre-flight* check on cluster health, not a guarantee about how many copies ultimately received the write; that distinction is a favourite interview follow-up. 3. **Apply locally.** The primary validates the request — mapping compliance, any `if_seq_no`/`if_primary_term` condition — then applies it: assigns the next **sequence number** from the shard's counter, stamps the current **primary term**, writes the document into the in-memory buffer, and appends the operation to the translog. With the default `index.translog.durability: request` the translog is fsynced here. 4. **Replicate.** The primary forwards the operation, carrying its assigned sequence number and term, **in parallel** to every shard copy listed in the in-sync set. Replicas do not re-derive anything; they apply the identical operation with the identical stamps, which is what keeps the copies byte-for-byte equivalent in ordering. 5. **Acknowledge.** Once every in-sync replica has responded successfully, the primary returns success to the coordinating node, which returns it to the client. So the default guarantee is strong: an acknowledged write is in the translog, fsynced, on the primary and on all in-sync copies. ## In-sync allocation IDs The cluster state records, per shard, the set of **in-sync allocation IDs** — the copies that are known to hold all acknowledged operations. This set, not the configured `number_of_replicas`, is what the primary replicates to and what the master is allowed to promote from. If a replica fails to apply an operation, times out, or its node drops, the primary does not fail the client's write. Instead it sends a shard-failed request to the master asking that the copy be removed from the in-sync set. Only after the master has applied that cluster-state update does the primary acknowledge. This ordering is the crucial safety property: a copy that missed an acknowledged operation must be marked stale **before** the client is told the write succeeded, so that copy can never be promoted to primary and silently lose the operation. If the primary itself cannot reach the master to record this, it steps down rather than continuing to acknowledge writes it cannot guarantee. ## Checkpoints Each copy tracks a **local checkpoint**: the highest sequence number below which it has processed every operation. The primary tracks the **global checkpoint**: the minimum local checkpoint across the in-sync copies, i.e. the point up to which every in-sync copy is complete. The global checkpoint is what makes recovery cheap — when a replica rejoins, the primary can often ship just the operations above that replica's checkpoint from its translog (operation-based recovery, kept feasible by retention leases) instead of copying whole segment files. ## Failure cases worth naming - **Replica node dies mid-write.** Primary marks it out of sync via the master; the write succeeds. Later the copy is re-allocated and recovered from the primary. - **Primary dies before acknowledging.** The client sees an error or a timeout, and the write's fate is genuinely ambiguous — it may or may not have reached replicas. This is why writes should be idempotent, ideally via a client-chosen document id. - **Primary is isolated but alive.** The master promotes a replica at a higher primary term. Operations from the old primary carry the old term and are rejected by the new one, which is the fencing mechanism. ## What this does *not* give you No cross-document or cross-shard atomicity. A bulk request touching five documents is five independent operations that can individually succeed or fail; there is no transaction spanning them. And an acknowledged write is durable but not yet *searchable* — that still waits for a refresh. ## Answer shape Walk the five steps, name the in-sync allocation set as the thing that is actually replicated to, and make the point that the master is asked to mark a failing replica stale before the client hears success. Adding local versus global checkpoints and their role in recovery is what turns a middle answer into a senior one.
- What does wait_for_active_shards actually guarantee?Only a pre-flight check: before starting, the primary verifies that at least that many copies of the shard are active in the cluster state, waiting up to a timeout if not. It does not guarantee the write reached that many copies — the primary still replicates to the in-sync set and may mark a copy stale mid-write. It reduces the chance of writing into a degraded shard, nothing more.
- Why must a failing replica be removed from the in-sync set before the client is told the write succeeded?Because the master may only promote copies that are in the in-sync set. If a replica missed an acknowledged operation but stayed listed as in-sync, promoting it later would silently lose that operation. Recording staleness first makes the acknowledgement honest: every copy still eligible for promotion has the write.
- How do the local and global checkpoints make replica recovery cheaper?Each copy's local checkpoint is the sequence number below which it has everything; the primary's global checkpoint is the minimum across in-sync copies. When a replica rejoins, the primary can replay just the operations above that replica's checkpoint from its translog — operation-based recovery — instead of copying whole segment files, provided a retention lease kept those operations available.
- A client gets a timeout from an index request. Was the document written?Unknown. The primary may have applied and replicated it before the response was lost, or it may have failed before applying. That is why writes should be idempotent: use a client-chosen document id so a retry overwrites the same document rather than creating a duplicate, and use conditional writes where the update is not naturally idempotent.
saying these in an interview costs you the question
- Says the primary replicates only after a refresh or flush
- Claims a failing replica causes the client's write to fail
- Thinks wait_for_active_shards guarantees how many copies got the write
- Believes replicas independently re-index and assign their own versions
- Treats a multi-document bulk request as an atomic transaction