Make search cheaper by reducing work, fan-out, and repeated computation without changing result meaning.

Search Tuning: Query Shape, Filters, Caches, Routing, Shard Count, Precomputation, and Aggregation Strategies

Build capacity and performance plans from measured workload dimensions and tail behavior rather than generic shard, heap, bulk, or hardware rules.

Intermediate → Advanced160–215 minutesSearch cost/fan-out experiment · Chapter 30 · Lesson 03Elasticsearch/Kibana 9.5.3 · OpenSearch/Dashboards 3.8.0 · free/local benchmark pathLast reviewed: September 2026

Learning outcomes

01

Reduce search cost by changing query shape and indexed data structures before reaching for thread-pool or cache-size knobs.

02

Separate scored query context from filters, request/query caches, filesystem cache, and application/session locality.

03

Quantify routing and shard fan-out benefits without creating hot spots or data skew.

04

Replace expensive repeated search-time computation with mapping/index-time precomputation when correctness permits.

05

Tune aggregations and pagination while preserving result correctness, relevance, and representative concurrency.

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.

Pinned performance baseline. Examples are reviewed against Elasticsearch/Kibana 9.5.3 (released 2026-09-03) and OpenSearch/OpenSearch Dashboards 3.8.0 (released 2026-08-04), using their bundled JVMs. The established local endpoints remain Elasticsearch at https://localhost:9200 and OpenSearch at https://localhost:9201 on the atlasmart-search Docker network. The generation environment did not execute live clusters, so this chapter never invents throughput, p95/p99, GC, disk, vector-recall, or cost results: numeric fields shown in report templates are deliberately blank/null until measured.

1. AtlasMart problem: a fast single query becomes a slow service

AtlasMart's “waterproof hiking” query looks fine in a developer console, but campaign traffic adds filters, facets, sorting, personalized routing, and concurrent users. Search capacity depends on the query's work per shard, the number of shards contacted, coordination/reduction cost, cache eligibility, fetch payload, and concurrency. Optimizing only the query clause with the largest Profile time can miss a fetch, aggregation, or fan-out bottleneck.

2. Query shape first: do less work

Use scored query context for conditions that should influence relevance and filter context for exact eligibility conditions such as tenant, availability, category, or time boundaries. Filter work can be reused/cached under the products' documented eligibility rules and avoids unnecessary scoring. Search only the fields that carry business meaning; broad multi-field or scripted queries can multiply work.

Lexical + exact filters with bounded response
GET atlasmart-perf-v1/_search
{
  "size": 10,
  "_source": ["sku","name","category","price"],
  "query": {
    "bool": {
      "must": [{"match":{"name":"waterproof hiking"}}],
      "filter": [
        {"term":{"tenant_id":"tenant-a"}},
        {"term":{"available":true}},
        {"range":{"price":{"lte":150}}}
      ]
    }
  }
}

3. Caches have eligibility and invalidation semantics

Do not optimize “cache hit rate” in isolation. The operating-system filesystem cache, node/query cache, and shard/index request cache store different things and have different invalidation/eligibility rules. In both products, request-level result caching is most useful for repeatable aggregation-style requests; dynamic requests involving relative now, profiling, or rapidly changing segments may not benefit. A warm-cache win is not a cold-start capacity guarantee.

Layer What it can save Capacity trap
filesystem/page cache disk reads for hot index files oversized heap or competing processes evict hot pages
query/filter cache repeated filter bitsets/results on eligible segments high churn/unique filters reduce reuse
request cache whole shard-level results for eligible requests frequent refresh invalidates entries; dynamic requests miss
application/session cache repeated business responses staleness and authorization/tenant leakage if keying is wrong

4. Routing can cut fan-out—and create hot shards

If AtlasMart indexes tenant-isolated documents with a stable routing key and searches with the same routing value, a request can touch fewer shards. That reduces per-request CPU/coordination. But routing values must be distributed enough to avoid one tenant or popular key making a shard disproportionately large or busy. Measure per-shard document count, indexing/search rate, CPU and disk; do not infer balance from total index size.

Compare unrouted and routed fan-out
GET atlasmart-perf-v1/_search
{"query":{"term":{"tenant_id":"tenant-a"}}}

GET atlasmart-perf-v1/_search?routing=tenant-a
{"query":{"term":{"tenant_id":"tenant-a"}}}

GET _cat/shards/atlasmart-perf-v1?v&h=index,shard,prirep,state,docs,store,node

5. Shard count sets the unit of parallelism and recovery

Each shard is a Lucene index. More primaries can increase parallelism and spread data, but every searched shard contributes coordination and per-shard work; too many tiny shards can deplete thread pools and fragment caches. Very large shards can lengthen recovery and reduce placement flexibility. There is no universal shard size that substitutes for measuring your workload, hardware, recovery objective, and growth pattern.

Wrong approach: choose “N GB per shard” as a law and resize the cluster until the number fits. Repair: benchmark candidate shard layouts with production data/query/write mix, then include a node-loss recovery test and future growth in the decision.

6. Precompute repeated business logic when it is stable

If every query computes the same normalized category, price band, recency class, or join-like derivation, consider materializing it during ingest or indexing. This trades write/storage cost for cheaper search. The change is correct only if the derived field has an explicit update rule; stale precomputed fields can create faster but wrong results.

7. Aggregation strategy is part of capacity

High-cardinality terms aggregations, large size/shard_size, scripts, nested traversals, and wide date histograms can increase memory and coordinating reductions. Use bounded user-facing facets, composite pagination when the use case requires enumerating many buckets, and pre-aggregated/rollup-style data when the business question does not require raw-document fidelity. Approximate metrics such as cardinality estimates must be accepted explicitly by the product requirement.

8. Pagination and fetch still count

Do not benchmark only query execution with size:0 if production fetches large _source payloads. Deep from/size can make every shard retain larger top-N sets; the Chapter 16 guidance remains: use stable search_after and PIT where appropriate, and keep PIT lifetimes bounded. Tail latency can be fetch-bound even after query clauses are optimized.

9. Diagnostic tools are not benchmark results

The Profile API is valuable for understanding relative query/aggregation costs but adds overhead, so do not quote Profile timings as production latency. Use Profile/Explain/slow logs/hot threads to form a hypothesis, then rerun the normal request in the fixed benchmark harness to prove or falsify it.

Capture server-side evidence before and after each run
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

10. AtlasMart controlled search 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.

  1. Run the same query mix at load level A and B with a fixed dataset/warmup.
  2. Record p50/p95/p99, throughput, searched shards, cache state, CPU/heap/GC/disk, errors/rejections, and relevance.
  3. Choose one change: move eligibility clauses to filters, add bounded routing, materialize a derived field, narrow fields/payload, or change an aggregation strategy.
  4. Do not simultaneously change shard count and query shape; that prevents causal attribution.
  5. Repeat the same runs and accept only if the intended SLO improves without a correctness/relevance or failure-recovery regression.
Benchmark report: fill only measured values
{
  "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}
}

11. Production judgment

Search tuning should reduce work while preserving result meaning. If a “faster” query changes tenant boundaries, facet totals, ranking judgments, or freshness, it is a defect. If a cache-dependent win vanishes after restart or during indexing, document that boundary in the capacity plan.

Check your understanding

  1. Why are filters often cheaper than scored clauses?
  2. Why can routing improve latency?
  3. What is the danger of custom routing?
  4. Why is Profile not a benchmark?
  5. What does Lesson 4 add?
Review the answers

1. They express eligibility without score computation and may reuse cacheable structures when eligible.

2. It can reduce the number of shards contacted, but only when indexing and searching use a compatible routing contract.

3. Skew/hot spots can concentrate data and traffic on a shard, trading fan-out reduction for imbalance.

4. Profiling adds diagnostic overhead and is intended to compare relative component cost, not measure normal request latency.

5. It maps the observed bottleneck to hardware/topology—CPU, heap, page cache, storage, network, tiers, and failure headroom—rather than changing software knobs blindly.

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

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.