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.

Intermediate115–145 minutesDistributed document/shard labElasticsearch 9.5.3 · OpenSearch 3.8.0Last reviewed: September 2026

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.

01

Trace an index request from receiving/coordinating node through routing to one primary shard and its replica group.

02

Explain what wait_for_active_shards checks before a write and what the response _shards block says afterward.

03

Trace query-then-fetch style scatter/gather across primary-or-replica shard copies.

04

Read _shards.total, successful, failed and partial-search controls instead of treating HTTP 200 as complete success.

05

Diagnose the different impact of an unavailable replica versus an unavailable primary in a disposable topology.

Version baseline reviewed 10 September 2026

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.

Execution and safety note

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.

REST · index one document and inspect write metadata
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.

REST · observe the active-shard precondition on a one-node lab
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.

What the failure proves

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.

REST · compare full fan-out and routed search
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?
REST · require complete shard execution for a critical query
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.

REST · product-neutral evidence once a supported two-node cluster exists
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=2 fails/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

  1. Why does one document write not execute on every primary shard?
  2. What does wait_for_active_shards=2 check, and what does it not guarantee?
  3. Why can a replica improve search capacity but also increase indexing cost?
  4. What should a client inspect even when a search returns HTTP 200?
  5. 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

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.