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.
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.
Trace the default query-then-fetch path from coordinating node to shard collectors and final fetch.
Explain can-match/pre-filter pruning and why it is an optimization rather than a correctness shortcut.
Distinguish local shard statistics from DFS global term statistics and identify the extra coordination cost.
Measure fan-out, skipped shards, query/fetch work and coordinator effects without treating one timing as the whole request.
Choose routing, shard topology or query changes from evidence instead of assuming every slow search is a Query DSL problem.
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.
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.
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"}]
}
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.
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
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"}
}
}
}
# 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.
# 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
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
- What does the query phase return before fetch?
- What does a skipped shard prove?
- Why can DFS change ranking?
- Why is DFS not a default tuning recommendation?
- 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
- Elastic search API — Current search request options, search type, pre-filter shard behavior and partial-results controls.
- Elastic search profiling — Profile operator trees, collectors and DFS profiling; also documents what profile does not measure.
- Elastic pagination — from/size result window, search_after, PIT consistency and PIT cleanup guidance.
- Elastic point in time API — Opening PITs, keep_alive, changing PIT IDs and retained-segment resource implications.
- Elastic async search submit — Long-running asynchronous search submission and response-size constraints.
- Elastic async search results — Polling, ownership/security and keep_alive behavior for async search.
- Elastic slow logs — Shard-level query/fetch slow logs and current query-logging guidance.
- OpenSearch pagination — from/size, search_after and PIT-backed pagination tradeoffs.
- OpenSearch PIT — Create/list/delete PIT APIs, security permissions and resource lifetime.
- OpenSearch profile API — Search component timings and explicit omissions such as network/queue/coordinator idle time.
- OpenSearch asynchronous search — Plugin endpoint, partial results and long-running search model.
- OpenSearch async settings — Maximum running time, concurrency, retention and wait timeout settings.
- OpenSearch logs — Request-level and shard-level search slow logs plus task-resource logging.
- OpenSearch tasks API — Task inspection and cancellation mechanisms.