Chapter 13 · Shard Sizing, Routing, Allocation, Awareness, and Cluster Topology

Primary/Replica Shard Count, Shard Overhead, Dataset Growth, Search Concurrency, and Recovery Time

Choose shard and replica counts from AtlasMart growth, query fan-out, recovery objectives and failure-domain constraints instead of copying a GB-per-shard rule of thumb.

Intermediate115–145 minutesShard topology & failure-domain labElasticsearch 9.5.3 · OpenSearch 3.8.0Last reviewed: September 2026

Learning outcomes

AtlasMart is growing from a small catalog into a multi-tenant search service. The dangerous shortcut is to pick a shard size from a blog post and treat it as an invariant. A shard is a Lucene index and a distributed scheduling/recovery unit: every extra shard has metadata, heap, file-handle, segment, search-coordination and recovery cost, while every oversized shard lengthens recovery and can limit parallelism.

01

Explain primary and replica shard roles without treating replicas as backups.

02

Model shard count from measured store size, growth, search concurrency and recovery objectives.

03

Relate query fan-out and per-shard work to p95/p99 latency rather than assuming more shards are faster.

04

Recognize both tiny-shard overhead and oversized-shard recovery risk.

05

Create a capacity decision record with explicit measurements and re-evaluation triggers.

Chapter baseline reviewed 11 September 2026

Examples target self-managed Elasticsearch 9.5.3 / Kibana 9.5.3 and OpenSearch 3.8.0 / OpenSearch Dashboards 3.8.0. The established AtlasMart endpoints remain https://localhost:9200 for Elasticsearch using its copied CA and https://localhost:9201 for the disposable OpenSearch demo certificate. Use each distribution's bundled JVM for this lab unless its current support matrix says otherwise. The initial Chapter 01 single-node containers are useful for API inspection but cannot demonstrate replica placement or zone resilience, so this chapter uses deterministic traces and an optional disposable multi-node topology. No moving latest tags and no production allocation changes are used.

Execution and safety note

Run mutating, destructive, security, lifecycle, snapshot, failure-injection, and load-test commands only in the disposable AtlasMart lab or an equivalently isolated environment. Verify the target cluster, index, tenant, credentials, and rollback path before execution; treat shown output as an expected invariant unless the lesson explicitly labels it as captured evidence.

1. A shard plan is a capacity model, not a storage slogan

A primary shard owns one partition of an index. A replica shard is another copy of that primary used for availability and, depending on request routing and workload, additional read capacity. Replicas do not replace snapshots because a bad delete, bad mapping migration or application bug can be replicated faithfully.

The useful question is not “how many GB should a shard be?” but “how much work can one shard absorb, how many shards can one node operate, and how long can one lost copy take to recover under the failure budget?” The answer changes with mapping shape, segment count, indexing rate, merge pressure, storage, cache residency, query complexity, aggregation cardinality and network throughput.

Decision input What to measure Why it changes the plan
Stored/indexed bytes store.size, primary store bytes, source/compression ratio Determines transfer/recovery volume and disk footprint.
Write rate docs/s, bytes/s, merge/indexing pressure A shard that is comfortable for reads may be too hot for writes.
Search fan-out shards touched per request and concurrent searches Coordination and per-shard task overhead grow with fan-out.
Recovery objective maximum acceptable degraded period after node/zone loss Constrains how large a lost shard copy can be.
Failure domain node/rack/zone count and replica policy Determines whether copies survive a correlated failure.
Growth horizon forecast plus uncertainty band Prevents a plan that is correct only on launch day.

2. Observe the topology before calculating anything

Dev Tools · inspect the current shard and node state
GET _cluster/health?level=shards
GET _cat/nodes?v&h=name,ip,node.role,heap.percent,ram.percent,cpu,disk.avail
GET _cat/shards/atlasmart-products*?v&h=index,shard,prirep,state,docs,store,node
GET atlasmart-products*/_stats/store,docs,indexing,search?level=shards
GET _nodes/stats/fs,indices,jvm,thread_pool

On the original single-node lab, a replica count greater than zero is intentionally unassignable. That yellow state proves only that there is nowhere legal to put the replica; it does not prove a data-loss event. For this chapter, keep the single-node environment for API syntax and use either the deterministic model below or a disposable multi-node cluster for placement experiments.

The CAT APIs are human-oriented snapshots. Automation should prefer structured JSON APIs where stable parsing matters. Capture topology, shard count, store size and request workload together; a shard count without workload context is not evidence.

3. Convert an RTO into a shard-recovery constraint

A rough lower-bound recovery model is copy_time ≈ shard_store_bytes / effective_recovery_bytes_per_second. “Effective” throughput is what the cluster can sustain while also serving production traffic; it is not the link speed printed on a NIC. Queueing, checksum work, source/destination disk, concurrent recoveries and throttles all matter.

Suppose AtlasMart measures 80 MiB/s of safe recovery throughput in a representative failure test and wants each lost shard copy to finish within 15 minutes. The arithmetic lower bound is roughly 72 GiB per copy. That is not a recommended shard size: it ignores queueing and competing recoveries, so the design needs safety margin and empirical verification.

Python · parameterized recovery envelope
from math import ceil

def recovery_floor_gib(effective_mib_s: float, minutes: float) -> float:
    return effective_mib_s * 60 * minutes / 1024

safe_mib_s = 80       # replace with measured sustained recovery rate
objective_min = 15    # AtlasMart example objective
print(round(recovery_floor_gib(safe_mib_s, objective_min), 1), "GiB lower-bound transfer envelope")

# Then test real shard sizes below this envelope under concurrent recovery + live traffic.
Do not turn the arithmetic into folklore

The equation is a planning bound, not a production guarantee. Measure recovery with realistic concurrent traffic and enough failed copies to expose queueing. A single 10 GiB test shard recovering on an idle cluster cannot validate a zone-loss RTO.

4. Search concurrency and fan-out pull in the other direction

Search is distributed. A query that targets an index normally executes shard-level work on every relevant shard copy and a coordinating node reduces the partial results. More primaries can provide parallelism and spread indexing, but they also create more shard tasks, more result reduction, more caches/segments and more cluster metadata.

A query touching 24 shards is not automatically slower than one touching 3, but it has a larger coordination surface. Measure p50/p95/p99, CPU, search thread-pool queue/rejections, heap and cache behavior while varying shard count with the same corpus. Keep data, query mix, concurrency and warmup identical.

Dev Tools · record fan-out and search pressure
GET /atlasmart-products/_search_shards
GET /_nodes/stats/thread_pool,indices?filter_path=nodes.*.thread_pool.search,nodes.*.indices.search
GET /_cat/thread_pool/search?v&h=node_name,name,active,queue,rejected,completed

5. Deliberately wrong approach: “one shard per X GB”

A fixed GB rule fails in both directions. Thousands of tiny shards can consume heap and file descriptors and create coordination overhead even if disk usage is low. One enormous shard can make recovery miss the RTO and reduce the scheduler’s ability to distribute work. The correct repair is to define measurable objectives and benchmark candidate topologies.

Also avoid the opposite myth that replicas always double throughput. Replica copies can serve searches, but the benefit depends on concurrency, cache locality, node saturation and request routing; replicas also multiply storage and recovery work.

Candidate Possible benefit Failure mode to test
Few larger primaries Lower fan-out and metadata overhead Long recovery; one hot shard can dominate.
More smaller primaries More placement/scheduling flexibility Tiny-shard overhead; high fan-out and merge concurrency.
More replicas Availability and potential read concurrency Storage/recovery multiplier; does not solve hot routing key.

6. AtlasMart lab: produce a shard decision record

Use a representative copy of AtlasMart data or a deterministic synthetic corpus. Do not extrapolate from document count alone: two documents with the same JSON size can create different index footprints because analyzers, doc values, stored fields, vectors and nested objects change storage and segment work.

Lab worksheet · record evidence, not guesses
Dataset:
  documents: <measured>
  primary_store_bytes: <measured>
  avg_ingest_docs_s / p95: <measured>
  query_mix + concurrency: <defined>
  candidate_primary_counts: [2, 4, 8]
  replicas: 1
  zones: 3

For each candidate:
  - p50/p95/p99 query latency after identical warmup
  - indexing throughput and rejection count
  - shard count per node and segment count
  - node CPU/heap/fs cache observations
  - single-node-loss recovery duration
  - zone-loss recovery duration or deterministic transfer model
  - remaining disk/CPU headroom

Reject any candidate that misses an objective even if average latency looks good.
  • Keep the workload fixed while comparing shard counts.
  • Record primary store bytes separately from replica-inclusive store.
  • Test recovery while traffic continues.
  • Document the next re-evaluation trigger: data growth, new vector field, new query class, new node type, or changed RTO.

Production judgment

Shard topology is a reversible decision only up to a point. You can create a new index and reindex, and some split/shrink operations exist under constraints, but every topology change consumes time and capacity. Leave enough headroom to migrate rather than running the cluster at the edge of disk or CPU.

On Elastic Cloud and Amazon OpenSearch Service, instance layouts, availability-zone configuration and platform limits may constrain what you can tune directly. Use the provider’s topology controls rather than assuming self-managed node attributes are exposed. Preserve snapshots independently of replica topology.

Check your understanding

  1. Why are replicas not backups?
  2. Why is “50 GB per shard” not a universal rule?
  3. What is search fan-out?
  4. What should validate a recovery-time calculation?
  5. When should the shard plan be revisited?
Review the answers

1. They copy the same logical index state, including accidental or malicious changes; recovery from logical corruption requires an independent backup such as snapshots.

2. Recovery throughput, query complexity, write rate, hardware, failure objectives and segment/mapping shape differ by workload.

3. The number of shard-level executions a request must coordinate; routing and index targeting can reduce it.

4. A representative failure test with realistic concurrent traffic, throttles, disk/network resources and multiple recovering copies.

5. When workload, corpus size, mapping, query mix, hardware, topology, RTO/RPO or managed-service constraints materially change.

Summary and next step

You now have a workload-driven way to reason about primary count, replicas, fan-out and recovery. Lesson 2 adds custom routing: a tool that can reduce fan-out but can also concentrate traffic and create hidden correctness assumptions.

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.