Chapter 02 · Documents, Indices, Data Streams, Shards, Replicas, Nodes, and Cluster Architecture
Distributed Write and Search Request Flows: Primary Shard, Replication, Scatter/Gather, and Partial Failures
Trace write replication and search scatter/gather across primary and replica shard copies, including active-shard and partial-result boundaries.
Learning outcomes
AtlasMart can now name every storage and control-plane component, but architecture becomes useful only when learners can follow a request. A document write does not broadcast to every shard. A search does not execute once against a monolithic index. This lesson traces the coordinating stage, primary stage, replica propagation, shard-level query work, global reduction and failure reporting.
Trace an index request from receiving/coordinating node through routing to one primary shard and its replica group.
Explain what wait_for_active_shards checks
before a write and what the response
_shards block says afterward.
Trace query-then-fetch style scatter/gather across primary-or-replica shard copies.
Read _shards.total, successful,
failed and partial-search controls instead of
treating HTTP 200 as complete success.
Diagnose the different impact of an unavailable replica versus an unavailable primary in a disposable topology.
This chapter continues the Chapter 01 lab baseline with
Elasticsearch 9.5.3 (released 3 September 2026)
and OpenSearch 3.8.0 (released 4 August 2026).
The existing local endpoints remain
https://localhost:9200 for Elasticsearch and
https://localhost:9201 for OpenSearch.
Elasticsearch requests use the Chapter 01 HTTP CA plus the
local elastic password; OpenSearch uses its local
demo TLS/admin path only for the disposable course lab.
Re-check versions, APIs and security defaults before reusing
these examples later.
The environment used to generate this chapter does not provide
Docker, Elasticsearch or OpenSearch, so commands were reviewed
against current official documentation but were not executed
here. Expected output is described by field shape and
invariant rather than fabricated as captured output. Failure
injection targets only disposable
atlasmart-* indices and the isolated Chapter 01
course containers; every destructive or allocation-changing
step includes an explicit reset path.
1. One document write targets one primary replication group
When a client indexes P-3001, the node that
receives the HTTP request acts as a coordinating node for that
operation. Routing resolves the document to one primary shard.
The request is forwarded internally to the current primary copy
for that shard. The primary validates and executes the
operation, assigns sequencing metadata, and propagates the
operation to active/in-sync replica copies according to the
product’s replication model.
This is a primary-backup data replication model. It does not mean the cluster has one global primary. An index with two primary shards has two primary replication groups; unrelated documents can be written through different primaries in parallel. The elected master/cluster-manager tracks routing/allocation state but is not the data primary for every write.
PUT /atlasmart-products-v2/_doc/P-3001?routing=TENANT-B{ "product_id": "P-3001", "name": "Atlas Split Keyboard", "category": "keyboards", "price": 129.00, "updated_at": "2026-09-10T17:00:00Z"}GET /_search_shards/atlasmart-products-v2?routing=TENANT-BGET /atlasmart-products-v2/_doc/P-3001?routing=TENANT-B
Correlate the routing result with the shard number shown by
_cat/shards. The write response’s
_index, _id, _seq_no,
_primary_term and _shards fields are
evidence from that operation. They are not a durable business
transaction receipt across every downstream system.
2. wait_for_active_shards is a precondition, not a
backup guarantee
Both current products expose
wait_for_active_shards. By default, document writes
can proceed once the primary is active. Requiring
2 or all asks the operation to wait
until that many shard copies are active before proceeding. This
reduces the chance of accepting a write with too little
redundancy, but the check occurs before the write; failures can
still happen while replication is executing. Inspect the
response afterward.
PUT /atlasmart-products-v2/_settings{"index": {"number_of_replicas": 1}}# One node cannot host the replica next to its primary.GET /_cluster/health/atlasmart-products-v2?prettyPUT /atlasmart-products-v2/_doc/P-3002?wait_for_active_shards=2&timeout=5s{ "product_id": "P-3002", "name": "Atlas Compact Keyboard", "category": "keyboards", "price": 59.00, "updated_at": "2026-09-10T17:05:00Z"}# RESET for the baseline lab.PUT /atlasmart-products-v2/_settings{"index": {"number_of_replicas": 0}}
On the one-node topology, the
wait_for_active_shards=2 write should wait and time
out/fail because a second shard copy cannot be active. The exact
error envelope is product/version dependent; do not fabricate
byte-identical output. After resetting replicas to zero, retry
with the normal precondition and verify the document.
It proves the requested active-copy precondition could not be satisfied on this topology. It does not prove that replicas are backups, that two active copies span independent failure domains, or that a later replica cannot fail during execution.
3. Search scatters to shard copies and gathers a global result
A search against atlasmart-products-v2 must cover
each relevant primary-shard group unless routing, index
selection or query planning prunes work. The coordinating node
selects an eligible copy from each target shard group, sends
shard-level search work, then merges/reduces results. For
top-hit searches, the common model is a query phase that
identifies competitive document IDs/scores followed by a fetch
phase that retrieves hit content. Aggregations add their own
shard-local and reduce semantics.
Replica shards can serve searches as well as primaries, so adding replicas can increase available search copies—if CPU, heap, filesystem cache, coordination and query patterns are balanced. It also increases indexing/storage/recovery work. More replicas do not automatically lower p99 latency.
GET /atlasmart-products-v2/_search{ "profile": false, "query": {"match": {"name": "keyboard"}}, "sort": [{"product_id": "asc"}]}GET /_search_shards/atlasmart-products-v2GET /_search_shards/atlasmart-products-v2?routing=TENANT-BGET /atlasmart-products-v2/_search?routing=TENANT-B{ "query": {"match": {"name": "keyboard"}}}
The unrestricted request addresses every primary-shard group in the index; the routed request narrows shard selection to the routing values supplied. That reduction is safe only when application semantics guarantee the desired documents were indexed with compatible routing. A missing routing value can turn a performance optimization into a correctness bug.
4. HTTP success can still contain shard-level failure information
Search responses expose a _shards summary.
Depending on product settings and request parameters, a search
can return partial results when some shards fail or time out.
Current Elasticsearch and OpenSearch both have controls around
allow_partial_search_results / default
partial-result behavior. A client that ignores
_shards.failed can display incomplete results as if
they were complete.
| Signal | Interpretation | Application question |
|---|---|---|
_shards.total |
Number of shard copies/groups targeted by the search operation as reported by the response. | Did the query touch more shards than expected? |
_shards.successful |
Shard operations that completed successfully. | Is this equal to the required completeness contract? |
_shards.failed |
Shard operations that failed. | Should the API fail closed, degrade with a warning, retry, or use fallback? |
timed_out |
Search exceeded its time budget in a way reported by the API. | Is returning partial/late data acceptable for this endpoint? |
GET /atlasmart-products-v2/_search?allow_partial_search_results=false{ "query": {"match": {"name": "keyboard"}}}
Do not globally set “no partial results” or “always allow partial results” by habit. A product-autocomplete endpoint might prefer a degraded but fast response; a compliance export or security investigation may require explicit failure if any shard is missing. The application contract decides.
5. Optional two-node proof: see a real replica copy without running four nodes at once
The resource-conscious Chapter 01 baseline runs one Elasticsearch node and one OpenSearch node simultaneously. To prove actual replica placement, stop the other product temporarily and expand one product into a two-node cluster using that product’s current official multi-node Docker procedure. Do not try to improvise node discovery/security parameters from this lesson if official 9.5.3/3.8.0 instructions differ—cluster bootstrap and certificate workflows are version-sensitive.
Once two eligible nodes are joined, set
number_of_replicas: 1 on
atlasmart-products-v2, then verify that each
primary and its replica are on different nodes. Stop the node
hosting one replica and inspect allocation/recovery. If the
primary remains active elsewhere, writes/searches may continue
according to shard availability; if the only active primary copy
for a shard disappears before promotion/recovery, that shard
becomes unavailable. Restore the node and observe convergence.
PUT /atlasmart-products-v2/_settings{"index": {"number_of_replicas": 1}}GET /_cluster/health/atlasmart-products-v2?wait_for_status=green&timeout=60sGET /_cat/shards/atlasmart-products-v2?v=true&h=index,shard,prirep,state,nodePUT /atlasmart-products-v2/_doc/P-3003?wait_for_active_shards=2{ "product_id": "P-3003", "name": "Atlas Ergonomic Keyboard", "category": "keyboards", "price": 139.00, "updated_at": "2026-09-10T17:10:00Z"}
Record the two node names and actual shard placements before failure injection. If your machine does not have enough memory to run a supported two-node product cluster, the mandatory learning path remains the one-node active-shard failure plus shard-allocation evidence; do not fake replica placement.
Verification checklist
- A routed write can be correlated with one target primary-shard group.
-
wait_for_active_shards=2fails/times out on the one-node + one-replica topology, then the lab is reset. -
Full search and routed search show different shard-selection
sets through
_search_shards. - Critical-query examples check shard completeness rather than trusting HTTP status alone.
- If the optional two-node lab is run, primary and replica copies are shown on different nodes with actual evidence.
Production judgment
Write latency includes coordination, primary execution and
replication; search latency includes shard-level work plus
coordination/reduction and fetch. Increasing shard or replica
counts changes both resource surfaces.
wait_for_active_shards encodes a precondition, not
a disaster-recovery strategy. Partial search behavior belongs in
endpoint-level correctness requirements. Before scaling a
cluster, measure fan-out, hit/aggregation sizes, replica
utilization, recovery time, disk/network bandwidth and
coordinating-node pressure rather than applying universal
shard/replica rules.
Check your understanding
- Why does one document write not execute on every primary shard?
-
What does
wait_for_active_shards=2check, and what does it not guarantee? - Why can a replica improve search capacity but also increase indexing cost?
- What should a client inspect even when a search returns HTTP 200?
- When is routed search a correctness risk?
Review the answers
1. Routing selects one primary-shard replication group for the document; other primary shards do not own that document.
2. It waits for the requested number of active copies before the operation proceeds; it does not guarantee those copies stay healthy, span independent failure domains, or replace snapshots.
3. Replicas can serve search requests, but each write must also be replicated and each copy consumes storage, cache and recovery bandwidth.
4. Shard success/failure counts, timeout indicators, failures and the actual result contract—not just transport status.
5. When the application supplies routing that does not cover all documents relevant to the query, the search can intentionally skip shards containing valid matches.
Summary and next step
Writes route to one primary replication group and propagate to replicas; searches scatter across relevant shard copies and gather results at the coordinating node. Active-shard preconditions and partial-search signals make failure visible but do not remove application policy. Next, trace one product all the way from request metadata through routing, indexing, refresh and returned search hit.
Authoritative references
- Elasticsearch reading and writing documents — Current primary-backup replication and coordinating/primary/replica stages.
- Elasticsearch node roles — Scatter/gather coordinating-node behavior.
-
Elasticsearch create-index active-shard semantics
— Current
wait_for_active_shardsframing. - OpenSearch document APIs — Current primary-backup document replication overview.
-
OpenSearch Index Document API
— Current
wait_for_active_shards, refresh and routing controls. - OpenSearch search shard routing — Current shard-copy selection, adaptive replica selection and custom routing behavior.