Chapter 13 · Shard Sizing, Routing, Allocation, Awareness, and Cluster Topology
Shard Rebalancing, Recovery, Throttling, Relocation, and Operational Headroom
Operate shard movement without turning recovery into a second outage: distinguish relocation, recovery and rebalance, measure progress, and preserve resource headroom.
Learning outcomes
AtlasMart replaces a data node after a disk fault. The cluster starts copying shards and latency rises. Operators must distinguish recovery (constructing a shard copy), relocation (moving a shard from one node to another) and rebalancing (allocator-driven movement toward a better distribution). All three compete with foreground traffic for disk, network, CPU and page cache.
Differentiate recovery, relocation and rebalancing from the shard states and APIs that expose them.
Monitor recovery bytes/files/stages and allocation throttling before changing concurrency.
Use node decommission filters safely and remove temporary constraints afterward.
Reason about recovery bandwidth/concurrency as a latency-versus-restoration tradeoff.
Define operational headroom that survives node/zone loss without relying on emergency tuning.
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. Read shard movement as a state machine
A shard copy can be unassigned, initializing, started or relocating. A relocation typically creates a target copy, transfers/reuses segment files and translog operations as needed, then hands over ownership and removes the old copy. A recovery may be local-store, peer-to-peer or snapshot-related depending on context.
Rebalancing is a policy decision that may trigger relocations; it is not a separate data-copy mechanism. During a node failure, allocation of missing copies takes precedence over cosmetic disk equality. Do not chase a visually equal shard count if latency and recovery objectives are already healthy.
GET _cluster/health?level=shards
GET _cat/shards?v&h=index,shard,prirep,state,node,relocating_node
GET _cat/recovery?v&active_only=true
GET _recovery?active_only=true&detailed=true
GET _cluster/pending_tasks
GET _cat/thread_pool?v&h=node_name,name,active,queue,rejected,completed
2. Recovery throttling protects foreground work
Both products expose cluster routing settings for concurrent
incoming/outgoing recoveries, with a general
node-concurrent-recoveries fallback. The allocator can report
THROTTLE when those limits are reached. There is
also recovery bandwidth control through recovery settings.
Defaults and managed-service overrides are version/platform
sensitive, so inspect rather than paste a tuning snippet.
Increasing concurrency can shorten the tail of a recovery queue on an idle cluster, but on a busy cluster it can raise search/index p99 by multiplying random I/O, network and cache churn. The safe target is the fastest recovery that still meets foreground SLOs.
GET _cluster/settings?include_defaults=true&flat_settings=true&filter_path=*.cluster.routing.allocation.node_concurrent_*,*.indices.recovery.*
GET _nodes/stats/fs,transport,indices,thread_pool,jvm
POST _cluster/allocation/explain
{
"index":"atlasmart-products-v2",
"shard":0,
"primary":false
}
A THROTTLE allocation decision often means the
system is intentionally limiting concurrent recovery. Raising
the limit without proving spare disk/network/CPU capacity can
turn degraded service into an outage.
3. Planned node removal: exclusion is safer than surprise
For self-managed maintenance, an allocation exclusion can drain shards from a node before shutdown. This is still real data movement: first verify remaining nodes have enough disk, failure-domain coverage and performance headroom. An exclusion that cannot be satisfied creates unassigned shards or endless relocation attempts.
Use a disposable node name in labs; never copy the command against an unrelated cluster. Remove the exclusion after maintenance so future allocations can use the node again.
PUT _cluster/settings
{
"persistent":{
"cluster.routing.allocation.exclude._name":"atlasmart-data-3"
}
}
GET _cat/shards?v&h=index,shard,prirep,state,node,relocating_node
GET _cluster/health?wait_for_no_relocating_shards=true&timeout=60s
# After the node is safely returned and verified:
PUT _cluster/settings
{
"persistent":{
"cluster.routing.allocation.exclude._name":null
}
}
4. Operational headroom is part of the topology
A six-node cluster that requires all six nodes to meet normal traffic has effectively zero node-failure headroom. Capacity plans must describe the degraded topology: one node down, one zone down, and recovery in progress. The target is not “green at any cost”; it is to keep the service within agreed latency/error objectives while restoring redundancy in the allowed time.
Reserve disk for relocation targets and merge amplification, CPU for extra shard/search work on survivors, network for copies, and heap/page cache for changed shard placement. A node that is 89% full under an Elasticsearch/OpenSearch 90% high-watermark default has almost no practical relocation room even though it has not crossed the threshold yet.
| Resource | Failure-time demand | Headroom evidence |
|---|---|---|
| Disk | Temporary source + target copies, merges, translog | Free bytes by node/tier/zone; watermark margin. |
| Network | Peer recovery plus client traffic | Recovery MB/s and transport saturation. |
| CPU | Checksum/decompression/index/search on survivors | p95/p99 CPU and search/index latency. |
| Heap/page cache | Different shard/segment working set | GC, cache churn, breaker and latency signals. |
| Thread pools | Recovery plus search/write queues | Queue/rejection counters during game-day. |
5. Deliberately wrong approach: maximize recovery speed
A common incident response is to raise every recovery limit until the cluster turns green faster. That optimizes one metric—time to redundancy—while potentially violating the user-facing objective. The mirror-image mistake is throttling recovery so heavily that the cluster stays under-replicated for hours and increases exposure to a second failure.
Repair the runbook by defining two simultaneous acceptance criteria: foreground p95/p99/error budget and redundancy-restoration RTO. Tune or scale only when measurements show which constraint is binding.
6. AtlasMart recovery experiment
Preferred path: run an optional disposable multi-node lab with small synthetic indices, stop one data-node container, and record recovery progress. If your machine cannot support it, use a deterministic trace with measured shard bytes and a parameterized recovery rate. Do not simulate disk pressure by filling the host disk.
Run at least two states: idle recovery and recovery under a fixed search/indexing workload. The lab is complete only when you can explain why the foreground latency distribution changes and which resource saturated first.
Scenario: one data node unavailable
Corpus: <bytes/docs measured>
Primaries/replicas: <record>
Zones/awareness: <record>
Foreground load: <queries/s + writes/s>
Capture every 10-30 s:
active recoveries + bytes remaining
p50/p95/p99 application latency
node CPU / disk / network
search + write queue/rejections
cluster health + unassigned count
Stop condition:
recovery RTO met AND foreground SLO not violated
OR abort if resource/latency safety bound is crossed.
Production judgment
Recovery tuning is temporary operational policy, not a substitute for enough nodes and disk. Keep baseline settings in configuration management, audit emergency overrides, and reset them after the incident. Managed services may hide or constrain recovery settings; the same principles still apply through provider metrics and scaling controls.
Check your understanding
- How does relocation differ from rebalancing?
- Why can faster recovery hurt users?
- What does a THROTTLE allocation decision mean?
- Why clear an exclusion after maintenance?
- What are the two key game-day objectives?
Review the answers
1. Relocation is the movement of a shard copy; rebalancing is an allocator decision that may cause relocations to improve distribution.
2. More concurrent copying can consume disk, network, CPU and cache resources needed by foreground search/indexing.
3. A legal allocation is temporarily limited by a throttling rule, often because concurrent recovery limits are reached.
4. A forgotten exclusion permanently removes capacity from the allocator and can block future legal placements.
5. User-facing latency/error objectives and the time allowed to restore redundancy.
Summary and next step
You can now observe and bound shard movement instead of treating green health as the only objective. Lesson 5 combines the chapter into a capacity/shard plan and a deterministic node-loss acceptance test.
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 CAT recovery API — Human-readable active and completed shard recovery status.
- OpenSearch CAT recovery — Recovery progress across OpenSearch shards.