Chapter 16 · Search Execution, Profiling, Slow Logs, Async/Search-After/PIT, and Pagination

Distributed Query/Fetch Phases, Can-Match, Shard Pruning, DFS Concepts, and Coordination Cost

Trace distributed search from shard selection through query, reduce and fetch, then reason about can-match pruning, DFS term statistics, coordinator work and the cost of fan-out.

Intermediate → Advanced115–150 minutesDistributed search & pagination labElasticsearch 9.5.3 · OpenSearch 3.8.0Last reviewed: September 2026

Learning outcomes

AtlasMart product search is fast when it targets one shard but develops a long tail after the catalog is split across many shards. The query itself has not become semantically more complex; the coordinator now has more shard work to schedule, merge and fetch. This lesson makes that distributed execution visible before anyone reaches for a random query rewrite.

01

Trace the default query-then-fetch path from coordinating node to shard collectors and final fetch.

02

Explain can-match/pre-filter pruning and why it is an optimization rather than a correctness shortcut.

03

Distinguish local shard statistics from DFS global term statistics and identify the extra coordination cost.

04

Measure fan-out, skipped shards, query/fetch work and coordinator effects without treating one timing as the whole request.

05

Choose routing, shard topology or query changes from evidence instead of assuming every slow search is a Query DSL problem.

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. AtlasMart keeps https://localhost:9200 for Elasticsearch with CA verification and https://localhost:9201 for the disposable OpenSearch demo certificate using OPENSEARCH_INITIAL_ADMIN_PASSWORD. The established containers are atlasmart-es and atlasmart-os. This chapter deliberately creates a disposable index with three primary shards and zero replicas to make shard fan-out observable on one local node; that is a teaching topology, not a production sizing recommendation. OpenSearch demo -k remains local-only; production must validate certificates. No moving latest tags are used.

Execution note

The generation environment does not run the AtlasMart containers, so latency, profile nanoseconds, slow-log lines, task IDs and PIT IDs are not fabricated. Expected outputs describe invariant fields and directions of change. Run the bounded lab locally and record your own p50/p95/p99, shard counts, profile trees and resource statistics before accepting a performance conclusion.

1. Search is a distributed top-k problem

A coordinating node accepts the client request and resolves the target index, alias or data stream into shard copies. During the query phase, each selected shard executes the query against its Lucene segments and returns candidate document identifiers plus score or sort values. The coordinator merges those shard-local candidate lists into a global top-k. During the fetch phase, the coordinator asks only the shards owning the winning documents for stored fields or _source, then assembles the final response.

This is why query cost and fetch cost can move independently. A cheap filter followed by fetching very large sources may be fetch-bound. A size-zero aggregation has essentially no hit fetch but can be query/aggregation-bound. Deep from/size makes every shard keep more candidates even though the client receives only one page.

Phase Shard work Coordinator work Useful evidence
Rewrite / pre-filter Rewrite query; sometimes prove shard cannot match. Resolve targets and optionally run pre-filter round. _shards.total, _shards.skipped, profile/rewrite behavior.
Query Score/filter/aggregate and keep local top candidates. Track shard responses and merge candidate sets. Profile query/collector tree, shard slow log, CPU/hot threads.
Reduce Return local candidates/aggregation partials. Merge top-k and reduce aggregations. Client took vs shard timings; coordinator CPU/queues.
Fetch Load selected stored fields/_source. Fan fetch requests to owning shards and assemble hits. Fetch slow logs, source size, returned fields, client bytes.

2. Can-match and shard pruning

When a request expands to many shards, Elasticsearch can run a pre-filter round based on query rewriting. A shard whose indexed range proves it cannot match a mandatory range predicate can be skipped. The pre_filter_shard_size request option controls when this extra round is enforced; current Elasticsearch also performs it automatically for conditions such as very large shard fan-out, read-only indices, or indexed primary sort cases.

OpenSearch performs equivalent distributed shard selection/rewrite optimizations but implementation details and defaults are version-sensitive. Do not code application correctness around a particular number of skipped shards. A shard being skipped means the engine proved it cannot contribute for that request; a shard not being skipped does not mean it will return a hit.

Observe shard fan-out and pruning
GET atlasmart-search-exec-v1/_search?pre_filter_shard_size=1
{
  "size": 5,
  "query": {
    "bool": {
      "filter": [
        {"range":{"updated_at":{"gte":"2026-09-01T03:00:00Z"}}},
        {"term":{"available":true}}
      ]
    }
  },
  "sort":[{"updated_at":"asc"},{"sku":"asc"}]
}
What the evidence proves

A response can expose total/successful/skipped/failed shard counts. It proves the request’s shard execution outcome, not that shard pruning is always beneficial; the extra pre-filter round itself has coordination cost.

3. query_then_fetch versus dfs_query_then_fetch

With the default query_then_fetch, relevance scoring uses term/document statistics local to each shard. On sufficiently representative shards, BM25 ranking is usually acceptable and avoids another coordination round. dfs_query_then_fetch adds a distributed-frequency-search phase that gathers global term statistics before the query phase, improving statistical consistency for some small/skewed shard cases at extra latency and coordination cost.

DFS is not a “make relevance better” switch. If the ranking defect comes from analyzers, field boosts, business signals or bad judgments, global term statistics do not repair it. Chapter 07’s rule still applies: measure ranking quality on judged queries.

Compare search types with the same query
GET atlasmart-search-exec-v1/_search?search_type=query_then_fetch
{
  "profile": true,
  "query":{"match":{"name":"wireless headset"}},
  "size":5
}

GET atlasmart-search-exec-v1/_search?search_type=dfs_query_then_fetch
{
  "profile": true,
  "query":{"match":{"name":"wireless headset"}},
  "size":5
}

On current Elasticsearch, DFS timing appears in profile output when DFS is used. Compare relevance order, not raw _score magnitude across different search modes as though score were a probability.

4. Lab: make coordination cost observable

Disposable Chapter 16 index — run separately against each product
PUT atlasmart-search-exec-v1
{
  "settings": {
    "number_of_shards": 3,
    "number_of_replicas": 0,
    "index.max_result_window": 100
  },
  "mappings": {
    "properties": {
      "sku":        {"type":"keyword"},
      "name":       {"type":"text","fields":{"raw":{"type":"keyword"}}},
      "category":   {"type":"keyword"},
      "price":      {"type":"double"},
      "available":  {"type":"boolean"},
      "updated_at": {"type":"date"},
      "popularity": {"type":"integer"}
    }
  }
}
Deterministic fixture generator (Python 3, writes NDJSON)
# Save as chapter16_fixture.py and redirect stdout to fixture.ndjson.
import json
from datetime import datetime, timedelta, timezone
base = datetime(2026, 9, 1, tzinfo=timezone.utc)
for i in range(240):
    sku = f"P-{1000+i:04d}"
    meta = {"index":{"_index":"atlasmart-search-exec-v1","_id":sku}}
    doc = {
      "sku": sku,
      "name": f"AtlasMart {'wireless' if i%3==0 else 'wired'} headset model {i:03d}",
      "category": ["audio","mobile","office"][i%3],
      "price": round(20 + (i%80)*1.25, 2),
      "available": i%5 != 0,
      "updated_at": (base + timedelta(minutes=i)).isoformat().replace('+00:00','Z'),
      "popularity": (i*17)%101
    }
    print(json.dumps(meta,separators=(',',':')))
    print(json.dumps(doc,separators=(',',':')))
# Then POST fixture.ndjson to /_bulk?refresh=true with Content-Type application/x-ndjson.
Index fixture and inspect topology
# Elasticsearch (verified CA)
curl --cacert atlasmart-es-http-ca -u elastic:$ELASTIC_PASSWORD   -H "Content-Type: application/x-ndjson"   --data-binary @fixture.ndjson "https://localhost:9200/_bulk?refresh=true"

# OpenSearch disposable demo TLS
curl -k -u admin:$OPENSEARCH_INITIAL_ADMIN_PASSWORD   -H "Content-Type: application/x-ndjson"   --data-binary @fixture.ndjson "https://localhost:9201/_bulk?refresh=true"

GET _cat/shards/atlasmart-search-exec-v1?v
GET atlasmart-search-exec-v1/_count
GET atlasmart-search-exec-v1/_search_shards
  • Confirm exactly 240 documents before performance comparisons.
  • Run the same query with and without a routing value only if the indexed documents used that same routing contract; arbitrary search routing can silently exclude documents.
  • Record _shards, client took/latency, profile tree and response bytes separately.
  • Repeat enough times to see a distribution; do not declare a win from one request.

5. Wrong approach: optimize only the query phase

Misleading fix

An engineer sees a cheap query profile, assumes “search is fine,” then increases client concurrency even though fetch returns full 200 KB product documents. Tail latency worsens because profile does not include every queue/network/coordinator cost and fetch remains dominant.

Repair by reducing fetched fields or source payload where semantics permit, measuring fetch slow logs/client bytes, and correlating coordinator/thread-pool evidence. Query profile is a microscope for parts of execution—not an end-to-end benchmark.

6. Production judgment

Shard fan-out is a topology decision as much as a query decision. More shards can increase parallelism but also coordinator work, per-shard overhead and tail-risk. Custom routing can reduce fan-out only when the application can preserve a correct routing contract and avoid hotspots. DFS can improve scoring consistency in narrow cases but spends another round trip.

For p95/p99 work, capture client end-to-end latency beside per-shard/query evidence. A request can have fast shard work and still be slow because of queueing, coordinator reduction, fetch payloads, remote-cluster/network latency, GC or backpressure.

Check your understanding

  1. What does the query phase return before fetch?
  2. What does a skipped shard prove?
  3. Why can DFS change ranking?
  4. Why is DFS not a default tuning recommendation?
  5. Why can profile look fast while the user is slow?
Review the answers

1. Shard-local candidate IDs with score/sort information and aggregation partials as applicable, not necessarily full source documents.

2. For that request, the engine proved the shard could not contribute; it does not establish a permanent routing rule.

3. It gathers global term/document statistics across shards before scoring instead of relying only on shard-local statistics.

4. It adds coordination and latency and only addresses a specific statistical issue, not general relevance problems.

5. Profile omits some end-to-end costs such as network, queueing and parts of coordination/fetch behavior.

Summary and next step

You can now trace query, reduce and fetch costs and recognize when shard selection or coordination dominates. Lesson 2 builds the diagnostic toolkit—profile, explain, slow logs, hot threads and tasks—without confusing any one tool with an end-to-end benchmark.

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.