skip to content

Why can two identical Elasticsearch searches rank results differently, and how does the preference parameter help?

level: seniorimportance: should knowfreq 40%

answer

  1. The two requests were not identical underneath
  2. Primary and replica are not byte-identical
  3. Deleted docs still count until merged away
  4. Pin the copy with a constant string
  5. Load-based copy selection is what you lose

basics

~20 s

Successive searches may be served by different copies of the same shard, and a primary and its replica hold different segment layouts and different numbers of not-yet-purged deleted documents, so their local statistics differ slightly. A constant preference value pins requests to the same copies.

solid answer

~50 s

Routing chooses the shard; the coordinating node then chooses which **copy** of that shard executes the search, balancing across the primary and its replicas. Copies are logically identical but not physically identical: they merge segments independently and can hold different numbers of deleted-but-not-yet-purged documents, so their local term statistics differ a little. Two identical searches served by different copies can therefore produce slightly different scores, and documents with near-equal scores swap places — the "bouncing results" problem, which is most visible while paging. Passing `preference` set to a constant string, such as a session or user id, makes the coordinating node hash that string to choose copies, so every request from that user hits the same copies and sees a stable ordering. Other values exist — `_local`, `_only_local`, `_shards:0,2` — but a constant custom string is the usual fix. The tradeoff: pinning bypasses the load-based copy selection and can concentrate traffic on one copy.

code

bash · 7 lines
bash
# stable ordering for one user's paging session, without pinning the whole cluster
GET /products/_search?preference=session-8f31c0
{
  "from": 20,
  "size": 10,
  "query": { "match": { "title": "tent" } }
}

go deeper

for a junior

Know that a search can be answered by either the primary or a replica of a shard, and that the choice is made per request by the node receiving the search.

for a middle

Explain why copies diverge physically — independent merges, deleted documents still counted until merged, unsynchronised refreshes — and how that produces slightly different scores.

for a senior

Diagnose bouncing results in production, apply a session-scoped preference, and articulate that you are trading adaptive load balancing for ordering stability.

for a principal

Decide where result stability is a product requirement versus a nice-to-have, and set the pattern for how session affinity is scoped so it never collapses the search tier onto one copy per shard.

## Shard choice versus copy choice Two separate decisions happen on the way to a search result. **Routing** decides which shards must be consulted. **Copy selection** decides, for each of those shards, whether the primary or one of its replicas actually runs the query. Copy selection is invisible in the response, which is why the symptom it causes looks mysterious. ## Why copies are not byte-identical A replica receives the same operations as its primary, so it holds the same documents. It does not hold them the same way: - **Segments merge independently.** Each copy runs its own merge scheduler against its own segment set, so at any instant the primary might have twelve segments and the replica seven. - **Deleted documents linger differently.** A delete or an update writes a marker; the old document remains in its segment until a merge rewrites that segment. A copy that has just merged has purged those documents; one that has not still counts them. - **Refreshes are not synchronised.** Each copy makes newly indexed documents searchable on its own refresh cycle, so for a fraction of a second one copy can see a document the other cannot. Because relevance scoring reads document counts and term document frequencies from the local index, and deleted documents still contribute to those counts until merged away, the two copies compute marginally different scores for the same document. ## The bouncing-results symptom On its own a difference in the fourth decimal place is harmless. It stops being harmless when documents tie or nearly tie — the classic case being a sort where many documents share the same value, or a filter-heavy query where scores cluster. Then the ordering of the tied group is decided by tiny differences, and it flips between requests. To a user paging through results this looks like documents disappearing from page 2 and reappearing on page 3, or a result list that reshuffles on every refresh. ## What preference does The `preference` query parameter tells the coordinating node how to choose copies: - **A custom string.** Any arbitrary value that does not start with `_`. The node hashes it to pick copies deterministically, so the same string always lands on the same copies for as long as the shard allocation is unchanged. Using a session id or user id gives each user a consistent view without pinning the whole cluster to one copy. - **`_local`.** Prefer copies on the coordinating node itself, falling back elsewhere. Cuts a network hop. - **`_only_local`.** Restrict to copies on the coordinating node, failing rather than going remote. - **`_only_nodes:<spec>`, `_prefer_nodes:<spec>`.** Constrain or bias copy selection to named nodes. - **`_shards:0,2`.** Restrict the search to specific shard numbers — a debugging tool, combinable with the others. ## Adaptive replica selection, and what pinning costs By default the coordinating node does not round-robin across copies. **Adaptive replica selection**, controlled by the cluster setting `cluster.routing.use_adaptive_replica_selection` and enabled by default in Elasticsearch 7.x and 8.x, ranks candidate copies using observed response times, the service time of previous requests, and the depth of the target node's search queue. A node that is garbage collecting, hosting a hot shard, or simply slower hardware receives less traffic automatically, which measurably improves tail latency in heterogeneous clusters. A constant `preference` overrides this. Every request carrying that value goes to the same copies regardless of how loaded they are. Pinning one user's session is harmless; pinning every request in the application to one literal string turns one copy of each shard into the entire search tier and leaves the replicas idle. This is the mistake to watch for. ## What preference does not fix A stable copy is not a stable snapshot. Both copies keep refreshing, so a document indexed between page 1 and page 2 will still change the result set even with `preference` set. Consistency across a paging session comes from taking a point-in-time view of the indices and paging within it; `preference` only removes the *copy-to-copy* variance on top of that. ## What interviewers listen for The strong answer separates shard routing from copy selection immediately, explains *why* two copies differ — merges and deletions, not "replication lag" in the relational sense — offers `preference` with a session-scoped value rather than a global constant, and names the load-balancing that pinning gives up.

  • Does setting preference guarantee two searches see exactly the same documents?
    No. It guarantees the same shard copies, which removes copy-to-copy score variance, but each copy keeps refreshing, so newly indexed or deleted documents still change the result set between requests. A stable view across a paging session requires a point-in-time snapshot of the indices, not a preference value.
  • What is wrong with hard-coding preference to a single literal string for the whole application?
    Every search then targets the same copy of every shard, so replicas sit idle while one copy of each shard absorbs the entire query load, and adaptive replica selection can no longer steer traffic away from a slow or garbage-collecting node. Scope the value per user or session instead.
  • How does adaptive replica selection decide which copy to use?
    It ranks the candidate copies using feedback from previous requests: the observed response time from each node, the service time those requests took, and the current depth of the node's search queue. Slower or busier nodes receive proportionally less traffic. It is a cluster-level dynamic setting and is on by default.

saying these in an interview costs you the question

  • Attributes the difference to replicas being stale copies
  • Says preference guarantees identical results between requests
  • Hard-codes one preference string for every user
  • Thinks the primary always serves searches
  • Cannot explain why deleted documents affect scores

context