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.
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.
Explain primary and replica shard roles without treating replicas as backups.
Model shard count from measured store size, growth, search concurrency and recovery objectives.
Relate query fan-out and per-shard work to p95/p99 latency rather than assuming more shards are faster.
Recognize both tiny-shard overhead and oversized-shard recovery risk.
Create a capacity decision record with explicit measurements and re-evaluation triggers.
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.
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
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.
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.
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.
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.
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
- Why are replicas not backups?
- Why is “50 GB per shard” not a universal rule?
- What is search fan-out?
- What should validate a recovery-time calculation?
- 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
- Elastic shard allocation settings — Current allocation, rebalance, disk-watermark and recovery controls.
- Elastic shard allocation awareness — Awareness and forced-awareness behavior across failure domains.
- Elastic allocation explain API — Explain why a shard is or is not allocatable.
- Elastic index recovery API — Recovery stage and byte/file progress.
- OpenSearch cluster settings — Current routing/allocation, awareness and recovery settings.
- OpenSearch cluster tuning — Shard allocation awareness and forced awareness concepts.
- OpenSearch CAT shards — Shard state and placement inspection.
- OpenSearch cluster allocation explain — Allocation diagnostics and decider evidence.
- Elastic size your shards — Guidance emphasizing workload testing rather than a universal shard-size number.
- OpenSearch shard indexing backpressure — Operational signals for shard-level indexing pressure.