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.

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

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.

01

Differentiate recovery, relocation and rebalancing from the shard states and APIs that expose them.

02

Monitor recovery bytes/files/stages and allocation throttling before changing concurrency.

03

Use node decommission filters safely and remove temporary constraints afterward.

04

Reason about recovery bandwidth/concurrency as a latency-versus-restoration tradeoff.

05

Define operational headroom that survives node/zone loss without relying on emergency tuning.

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. 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.

Dev Tools · movement dashboard without a GUI
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.

Inspect recovery controls and live evidence
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
}
Do not tune by symptom

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.

Pattern · drain and reset a disposable node
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.

Experiment record
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

  1. How does relocation differ from rebalancing?
  2. Why can faster recovery hurt users?
  3. What does a THROTTLE allocation decision mean?
  4. Why clear an exclusion after maintenance?
  5. 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

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.