Chapter 13 · Shard Sizing, Routing, Allocation, Awareness, and Cluster Topology
Model a Capacity/Shard Plan and Simulate Node Loss Without Exceeding Recovery or Latency Objectives
Turn AtlasMart workload assumptions into an explicit shard/capacity plan, exercise deterministic node-loss scenarios, and reject designs that cannot meet recovery or latency objectives.
Learning outcomes
The platform review now asks a concrete question: “Can AtlasMart survive one node or one availability-zone loss without missing the recovery objective or breaching search latency?” This lesson converts topology folklore into a small capacity model and an acceptance test. The model is deliberately conservative and parameterized; the numbers become valid only after you replace assumptions with measurements.
Turn dataset growth, ingest/search concurrency, node resources and failure objectives into candidate shard plans.
Model node and zone loss without claiming the model is a benchmark.
Define objective pass/fail thresholds for recovery, latency and disk headroom.
Validate routing skew and allocation legality alongside aggregate capacity.
Produce a change record that states assumptions, evidence, rollback 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. Start from service objectives and measured workload
Capacity planning should begin with what the service must do, not with how many nodes are currently available. Define a growth horizon, expected indexed bytes, writes/s, query mix/concurrency, p95/p99 targets, recovery RTO, tolerated availability state during a zone loss, and snapshot RPO. The last item is intentionally separate: replicas help availability; snapshots handle independent recovery.
Use ranges for uncertain inputs. A plan with 40% annual growth should also show what happens at a higher band, because reindexing under emergency disk pressure is the worst time to discover the forecast was optimistic.
| AtlasMart planning field | Example placeholder | Replace with |
|---|---|---|
| Primary store today | 600 GiB |
Measured primary-store bytes. |
| 12-month high case | 1.2 TiB |
Forecast + uncertainty. |
| Peak writes | 8k docs/s |
Observed p95/p99 interval. |
| Peak searches | 500 qps |
Query-class mix and concurrency. |
| Search objective | p99 < 350 ms |
Application SLO by query class. |
| Single-node recovery RTO | < 30 min |
Business/operations objective. |
| Zone-loss policy | Serve within SLO; replicas may be yellow under forced awareness | Explicit degraded-mode contract. |
2. Build a transparent model—not a magic calculator
The script below intentionally avoids pretending to know Elasticsearch/OpenSearch throughput. It takes measured values and calculates only simple envelopes: average primary bytes, replica-inclusive storage, a lower-bound single-shard copy time, and post-failure disk pressure. It cannot predict p99 query latency; that still requires a workload test.
A useful model exposes assumptions so reviewers can challenge them. Hidden constants are where “best practice” turns into folklore.
from dataclasses import dataclass
from math import ceil
@dataclass
class Plan:
primary_store_gib: float
primaries: int
replicas: int
data_nodes: int
node_disk_gib: float
safe_recovery_mib_s: float
@property
def avg_primary_gib(self):
return self.primary_store_gib / self.primaries
@property
def total_copies_gib(self):
return self.primary_store_gib * (1 + self.replicas)
@property
def avg_disk_used_pct(self):
return 100 * self.total_copies_gib / (self.data_nodes * self.node_disk_gib)
@property
def one_shard_copy_minutes_floor(self):
mib = self.avg_primary_gib * 1024
return mib / self.safe_recovery_mib_s / 60
for p in [
Plan(1200, 6, 1, 6, 1000, 80),
Plan(1200, 12, 1, 6, 1000, 80),
Plan(1200, 24, 1, 6, 1000, 80),
]:
print({
"primaries": p.primaries,
"avg_primary_gib": round(p.avg_primary_gib,1),
"avg_disk_used_pct": round(p.avg_disk_used_pct,1),
"one_shard_copy_min_floor": round(p.one_shard_copy_minutes_floor,1),
})
# Replace every value with measurements. This is not a benchmark or sizing recommendation.
3. Add failure-domain math before performance testing
With six equal data nodes across three zones and one replica, losing one zone removes roughly one third of data-node capacity at once. The survivors must serve traffic and reconstruct missing copies if the policy allows it. Average steady-state disk utilization can therefore look comfortable while zone-loss utilization becomes unsafe.
Compute a conservative survivor capacity ratio before running the cluster. If ordinary awareness will rebuild replicas into surviving zones, budget the temporary copy amplification. If forced awareness will leave replicas unassigned, budget the extra search/write load but not immediate full replica restoration until the missing zone returns.
def survivor_disk_pct(total_copy_gib, node_disk_gib, total_nodes, lost_nodes):
survivors = total_nodes - lost_nodes
return 100 * total_copy_gib / (survivors * node_disk_gib)
print("steady:", survivor_disk_pct(2400, 1000, 6, 0))
print("one node lost:", survivor_disk_pct(2400, 1000, 6, 1))
print("one 2-node zone lost:", survivor_disk_pct(2400, 1000, 6, 2))
# This assumes all copies are eventually present on survivors.
# Forced awareness may intentionally leave some replicas unassigned instead.
Equal-capacity arithmetic ignores shard indivisibility, allocation constraints, hot routing keys, merge/translog headroom and per-node workload. Use it to reject impossible plans early, not to certify a plan as safe.
4. Define the acceptance matrix before the game-day
A failure test without pass/fail criteria is a demo. Define the workload, failure, expected health state, maximum recovery time and application SLO before stopping anything. Keep failures isolated to disposable containers or use a deterministic trace if the local machine cannot host the topology.
| Scenario | Expected platform state | Pass evidence |
|---|---|---|
| No failure | All required copies started | Baseline p50/p95/p99 + no sustained rejection. |
| One data node lost | Missing copies recover/relocate legally | RTO met; p99 and error budget stay within degraded target. |
| One zone lost, awareness only | Copies may rebuild in surviving zones | Survivor disk/CPU remain below safety bound; recovery finishes. |
| One zone lost, forced awareness | Some replicas may remain unassigned/yellow | Service remains within SLO; reason for each unassigned replica is expected. |
| Skewed routing load | One key intentionally dominates | No single shard exceeds defined queue/latency/headroom limit. |
5. Execute with observable APIs and application metrics
GET _cluster/health?level=shards
GET _cat/nodes?v&h=name,node.role,cpu,heap.percent,ram.percent,disk.avail
GET _cat/shards?v&h=index,shard,prirep,state,docs,store,node,relocating_node
GET _cat/recovery?v&active_only=true
GET _cat/allocation?v
GET _nodes/stats/fs,indices,transport,thread_pool,jvm
POST _cluster/allocation/explain
{}
# Pair this with application p50/p95/p99, error rate, QPS/write rate and test timestamps.
Do not use _cluster/health alone. Green health can
coexist with terrible latency, and yellow can be an intentional
forced-awareness state. The acceptance artifact must pair
cluster state with user-facing workload metrics and explain any
shard that is unassigned or throttled.
6. Deliberately wrong approach: certify from averages
Average shard size, average CPU and average latency hide the failure domains that matter. One hot routing key, one slow disk, one oversized shard or one availability zone can determine the incident. p95/p99 and per-shard/per-node distributions are more useful than a single mean.
Another wrong pattern is to test node loss on an idle cluster. Recovery that meets RTO only when traffic is zero is not evidence for the production objective. Replay or generate a fixed representative workload during the experiment.
7. AtlasMart final shard-plan record
Write the plan as a versioned decision record so future chapters can reason about the same topology. Include: server versions, node roles/resources, zones, primaries/replicas, routing policy, lifecycle/tier rules, current and forecast primary bytes, workload classes, recovery settings, disk safety margin, snapshots, failure objectives, benchmark harness and the exact date of the last game-day.
Add re-evaluation triggers rather than pretending the decision is permanent. Examples: 25% primary-store growth, a new vector field, a high-cardinality aggregation workload, a new tenant that changes routing skew, new instance/storage type, changed zone count, or changed RTO/SLO.
- No universal shard-size threshold appears in the decision record.
- Every tuning override has an owner and reset condition.
- Replica topology is documented separately from snapshot/RPO policy.
- Managed-service restrictions are listed explicitly.
- The next migration path—reindex, split/shrink where applicable, or new index generation—is known before capacity becomes critical.
Production judgment
A shard plan is accepted evidence, not a static best practice. Keep the model in source control, keep the workload harness reproducible, and rerun the failure matrix after meaningful version/topology/workload changes. The cheapest cluster that passes steady-state benchmarks but fails a zone-loss test is not production-ready.
Check your understanding
- Why can a capacity model reject but not certify a topology?
- Why test with foreground traffic during recovery?
- What does forced awareness change in a zone-loss model?
- Which metric should accompany cluster health?
- What makes the shard plan maintainable?
Review the answers
1. Simple arithmetic can show impossible disk/recovery envelopes, but it cannot predict real query latency, cache behavior, merge pressure or allocator interactions.
2. Recovery competes for the same disk/network/CPU/cache resources as user requests; idle recovery is not representative.
3. It can intentionally leave replicas unassigned instead of rebuilding all copies into surviving zones, changing both disk demand and availability state.
4. Application latency/error/QPS plus per-node/per-shard resource and recovery evidence; health color alone is insufficient.
5. Versioned assumptions, reproducible measurements, explicit pass/fail objectives, rollback/migration options and re-evaluation triggers.
Summary and next step
You can now design shard topology from capacity, fan-out, routing distribution and failure objectives, explain allocator decisions, and validate node/zone loss without relying on folklore. Chapter 14 goes inside each shard to Lucene segments, refresh, merge, translog, flush and storage behavior—the mechanisms that make shard size and recovery cost tangible.
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 production guidance: size shards — Workload-driven shard sizing and benchmark guidance.
- Elastic CAT allocation — Per-node shard count and disk allocation view.
- OpenSearch CAT allocation — Per-node disk and shard allocation evidence.
- OpenSearch performance analyzer — Optional operational metrics for self-managed performance analysis.