Map measured bottlenecks to CPU, memory, storage, network, tiers, topology, and failure headroom.
Hardware/Topology: CPU, Heap, Page Cache, SSD, Network, Data Tiers, Hot/Warm/Cold, and Failure Headroom
Build capacity and performance plans from measured workload dimensions and tail behavior rather than generic shard, heap, bulk, or hardware rules.
Learning outcomes
Map CPU, heap, filesystem/page cache, storage, and network bottlenecks to the search/indexing mechanisms they constrain.
Use Elasticsearch automatic heap sizing and OpenSearch half-memory guidance as starting boundaries—not universal performance targets.
Plan SSD/local-vs-remote storage and data tiers from measured latency, merge, recovery, and cost behavior.
Reserve failure headroom so a surviving cluster can serve traffic while rebuilding lost shard copies.
Compare scale-up, scale-out, tiering, and workload separation as distinct topology choices with measurable tradeoffs.
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. AtlasMart problem: the cluster is “only 60% utilized” but p99 is failing
A node-level average can hide the actual bottleneck. One hot shard can consume a CPU core, a merge can saturate disk bandwidth, a large heap can starve filesystem cache, or cross-node fetch/recovery can fill the network. Capacity engineering maps each symptom to the resource used by that operation instead of treating RAM or CPU percentage as a single health score.
2. CPU: shard work, coordination, scripts, parsing, vectors
Search execution, scoring, aggregations, indexing analysis, ingest processors, compression, and some vector operations consume CPU. Because shard-level search work executes concurrently across requests, oversharding can amplify scheduling/thread-pool pressure. CPU saturation is meaningful only alongside achieved throughput and queue/rejection/tail behavior: a well-utilized CPU can be healthy at target load, while 70% average with bursty single-core hot spots can still violate p99.
3. Heap and page cache are a shared memory budget
Elasticsearch 9.5.3 automatically sizes heap by node
roles/available memory and recommends the default for most
production environments; if overridden, 50% is an upper-bound
guideline, not a target, and smaller can perform better by
leaving more filesystem cache. OpenSearch documentation uses
roughly half available memory as a starting point for
Xms=Xmx. In both products, leave memory for the OS,
native/off-heap structures, and page/filesystem cache.
| Memory consumer | Examples | Failure signal |
|---|---|---|
| JVM heap | cluster/index metadata, caches, aggregation/request state | GC pressure, breakers, long pauses, OOM risk |
| OS page/filesystem cache | Lucene segment/doc-values/index files | major faults, disk reads, cold-search latency |
| native/off-heap | network buffers, mmap/native vector libraries depending on engine | RSS growth beyond Xmx, native allocation pressure |
| other processes | agents, shippers, benchmark driver if co-located | stolen CPU/RAM and distorted benchmark results |
4. Storage: latency, throughput, IOPS, and merge/recovery concurrency
Fast local SSD often helps because indexing writes, merges, translog/fsync, cache misses, snapshots/restores, and shard recovery all use storage. But “SSD” is not a capacity number: measure sequential/random throughput, latency under concurrent merge/search/recovery, and actual cloud volume limits. Remote storage can be viable in supported architectures, but benchmark the exact path and failure semantics rather than transferring local-disk results.
5. Network becomes a search and recovery resource
Bulk clients send source bytes, coordinating nodes scatter searches and gather hits/aggregation reductions, replicas receive writes, and shard recovery transfers segment files. Cross-zone/region traffic also adds latency and cost. Record network bytes/packets and retransmission/error evidence during the same p99 window. A topology that passes normal load but saturates the network during node recovery has insufficient failure headroom.
6. Data tiers express access and cost intent
Hot/warm/cold-style tiers are useful when retention and query frequency change over time, but Elastic data tiers/ILM and OpenSearch ISM/allocation mechanisms are not interchangeable APIs. The capacity question is the same: what query latency, recovery objective, storage cost, and write activity must each age band support? Move only workloads whose SLO fits the target tier.
| Tier intent | Typical pressure | Capacity question |
|---|---|---|
| hot / actively written | CPU, fast storage, indexing buffers, merge | can it ingest + search + survive a failure at peak? |
| warm / mostly read | page cache and scan/aggregation cost | does slower/denser storage meet query SLO? |
| cold / infrequent | restore/remote access latency, cache misses | is access latency acceptable and recovery tested? |
7. Failure headroom is capacity, not waste
When a node fails, the surviving nodes must serve normal traffic and recover lost shard copies. Therefore a cluster sized to 100% of normal-state throughput has no room to heal. Define the failure domain—node, rack/zone, or service allocation—and test the degraded topology. Record client p99, indexing lag, rejections, recovery throughput, and time-to-green/recovered redundancy.
8. Scale up, scale out, or separate workloads?
| Choice | Can help when | Watch for |
|---|---|---|
| scale up | single-shard CPU, memory working set, local storage performance | larger failure unit; heap/page-cache balance; vertical limits |
| scale out | parallel shards, aggregate throughput, recovery placement | more shard/coordination overhead and network traffic |
| change shard layout | oversharding/undersharding is causal | reindex/rollover complexity and recovery size |
| separate workloads | indexing and search steal resources from each other | extra operational/storage cost and replication/remote-store semantics |
| tier data | older data has looser SLO and lower write activity | cross-tier queries and recovery behavior |
9. AtlasMart topology experiment
All Chapter 30 mandatory exercises preserve the course's
established free/local security boundary: Elasticsearch 9.5.3 at
https://localhost:9200 authenticated with
ELASTIC_PASSWORD and the copied CA
atlasmart-es-http-ca; OpenSearch 3.8.0 at
https://localhost:9201 authenticated with
OPENSEARCH_INITIAL_ADMIN_PASSWORD. The shared
Docker network remains atlasmart-search.
OpenSearch's -k examples are for the disposable
demo certificate only and are not production TLS guidance. The
default local fixture uses one primary and zero replicas because
it is a workstation lab; any availability or node-failure
conclusion must be tested in a disposable multi-node variant
rather than inferred from this single-node setup.
GET _cluster/health
GET _cat/nodes?v&h=name,cpu,heap.percent,ram.percent,disk.used_percent,node.role,master
GET _cat/thread_pool/search,write?v&h=node_name,name,active,queue,rejected,completed
GET _nodes/stats/jvm,process,os,fs,indices,thread_pool,indexing_pressure
GET atlasmart-perf-v1/_stats?level=shards
GET _cat/recovery/atlasmart-perf-v1?v
GET atlasmart-perf-v1/_segments
nodes:
data_nodes: <count>
cpu_per_node: <cores>
ram_per_node_gib: <GiB>
heap_policy: <automatic or Xms=Xmx>
storage: <local-ssd / cloud-volume / remote-store>
network: <Gbps / cross-zone>
shards:
primaries: <count>
replicas: <count>
working_set:
indexed_bytes: <bytes>
hot_bytes_estimate: <bytes>
failure_budget:
tolerated_node_failures: <count>
target_rto_s: <seconds>
peak_traffic_during_recovery: <percent-of-normal-or-target>
measurement:
p99_ms_normal: null
p99_ms_degraded: null
recovery_s: null
max_disk_pct: null
rejection_rate: null
If the mandatory workstation lab cannot create a meaningful multi-node failure domain, perform the topology decision with the deterministic template and use a bounded post-load recovery test. The optional free/local multi-node Docker variant may stop one data node only after snapshots/fixture replay are available and only if the cluster has enough eligible copies to remain correct.
{
"run_id": "atlasmart-perf-YYYYMMDD-NN",
"platform": {"product": null, "version": null, "topology": null},
"dataset": {"documents": null, "source_bytes": null, "indexed_bytes": null, "warm_state": null},
"offered_load": {"search_qps": null, "write_ops_s": null, "clients": null},
"results": {
"achieved_search_qps": null,
"achieved_write_ops_s": null,
"latency_ms": {"p50": null, "p95": null, "p99": null},
"indexing_lag_ms": {"p95": null, "p99": null},
"errors": null,
"rejections": null,
"relevance": {"ndcg_at_10": null, "vector_recall_at_10": null},
"recovery_seconds": null
},
"resources": {"cpu": null, "heap": null, "gc": null, "disk_io": null, "network": null},
"cost": {"currency": null, "window_cost": null, "cost_per_1k_queries": null},
"notes": {"warmup": null, "cache_state": null, "one_change_from_baseline": null}
}
10. Production judgment
Hardware choices are conclusions from a workload, not defaults. Rebenchmark when storage class, CPU generation, memory limit, vector engine, node role, shard layout, data tier, managed-service instance family, or network topology changes. Cost is part of the comparison: a topology that is 10% faster but 2× the cost may or may not be justified depending on SLO and failure risk.
Check your understanding
- Why can a larger heap make search slower?
- Why must storage be tested during recovery, not only steady search?
- What is failure headroom?
- Why is scale-out not automatically better?
- What does Lesson 5 prove?
Review the answers
1. It leaves less RAM for the filesystem/page cache and can increase GC pause costs; heap and OS cache share the physical-memory budget.
2. Recovery and merges can compete with serving traffic for I/O, exposing a bottleneck that normal-state tests miss.
3. Reserved capacity that lets surviving nodes serve the workload while rebuilding/recovering after a failure.
4. It adds nodes and parallelism but also shard, coordination, network, placement, and operational overhead.
5. Whether the complete production-like workload has a safe operating region, a measurable saturation knee, acceptable quality/freshness, and acceptable recovery/cost.
Summary and next step
Preserve the evidence, assumptions, version boundaries, and safety checks established in this lesson. Carry them into the next lesson—or, at the end of the capstone, into the production runbook—rather than treating this lesson as an isolated recipe.
References and current-version checks
- Elastic Stack 9.5.3 release
- Elasticsearch performance optimizations
- Elasticsearch tune for indexing speed
- Elasticsearch tune for search speed
- Elasticsearch size your shards
- Elasticsearch JVM settings
- Elasticsearch resilience guidance
- Elasticsearch bulk API
- Elasticsearch refresh parameter
- Rally documentation
- OpenSearch 3.8 version history
- OpenSearch indexing performance tuning
- OpenSearch Bulk API
- OpenSearch Refresh Index API
- OpenSearch search shard routing
- OpenSearch index settings
- OpenSearch shard indexing backpressure
- OpenSearch Benchmark quickstart
- OpenSearch Performance Analyzer metrics
- OpenSearch vector search performance tuning