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.

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

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.

01

Turn dataset growth, ingest/search concurrency, node resources and failure objectives into candidate shard plans.

02

Model node and zone loss without claiming the model is a benchmark.

03

Define objective pass/fail thresholds for recovery, latency and disk headroom.

04

Validate routing skew and allocation legality alongside aggregate capacity.

05

Produce a change record that states assumptions, evidence, rollback 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. 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.

Python · AtlasMart shard-plan envelope
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.

Python · simple survivor-capacity check
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.
Why this is only a guardrail

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

Evidence bundle to capture
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

  1. Why can a capacity model reject but not certify a topology?
  2. Why test with foreground traffic during recovery?
  3. What does forced awareness change in a zone-loss model?
  4. Which metric should accompany cluster health?
  5. 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

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.