skip to content

What are the two phases of a distributed SolrCloud query across shards?

level: middleimportance: should knowfreq 45%

answer

  1. The first round ships almost no data
  2. Each shard returns more than the client asked for
  3. Only the winners are fetched in full
  4. One parameter collapses it to a single trip
  5. Explains why deep offsets hurt

basics

~20 s

First, the receiving node asks one replica of every shard for the top matching document ids and sort values, and merges them into a global top-N. Second, it fetches the stored fields and highlighting for just those winning ids from the shards that hold them.

solid answer

~60 s

Any node can receive the query; it becomes the **aggregator** for that request. **Phase one** fans out to one replica of each shard, asking only for what it needs to rank: the document id, the score or sort field values, for `start + rows` documents per shard. Each shard must return `start + rows`, because the aggregator cannot know in advance how many of the global winners live on any one shard. The aggregator merges those lists and picks the global top-N. **Phase two** goes back to just the shards that own the winning ids and asks for the actual stored fields, plus highlighting or other per-document work. Two phases exist to avoid shipping full documents that lose the merge. Setting `distrib.singlePass=true` collapses it into one round trip, which pays off when `rows` is small and the field list is short. The design also explains why deep `start` values are expensive — each shard must produce `start + rows` results — and why `cursorMark` is the answer for deep paging.

code

bash · 2 lines
bash
# Per-phase, per-shard timings for a distributed request
curl "http://localhost:8983/solr/products/select?q=laptop&rows=10&debug=track"

go deeper

for a junior

Recall that any node can receive the query, fans it out to the shards, and merges the results into one response the client sees.

for a middle

Explain both phases concretely — ids and sort values first for start+rows per shard, stored fields for the winners second — and why the split saves work.

for a senior

Use it diagnostically: read debug=track timings, spot the slow shard, know when singlePass helps, and replace deep start offsets with cursorMark before they melt the cluster.

for a principal

Own the consequences at scale — shard count versus tail latency, whether inconsistent distributed scoring is acceptable for your relevance goals, and the paging contract the API exposes to clients.

## The aggregator A SolrCloud query can be sent to any node in the cluster. That node becomes the **aggregator** (sometimes called the coordinator) for the request: it decides which shards must be consulted, picks one replica per shard, issues the shard requests, merges the responses and returns a single result to the client. Nothing about the aggregator is special or preassigned — the role lasts exactly one request. Which shards are consulted depends on routing. Normally it is all of them; with a `_route_` parameter under the compositeId router, only the shard owning that routing prefix. Which *replica* of each shard is picked can be steered with `shards.preference` (for example preferring a same-node replica, or PULL replicas for analytics traffic). ## Phase one: find the winners The aggregator sends each shard a query asking for the minimum needed to rank: the uniqueKey, the score (or the values of the sort fields), and `start + rows` results. It deliberately does *not* ask for stored fields. The `start + rows` part is the important subtlety. If the client asked for `start=0&rows=10`, every shard must return its own top 10, because in the worst case all ten global winners live on a single shard. The aggregator merges the per-shard lists by score or sort key and takes the global top 10. ## Phase two: fetch the documents Now the aggregator knows exactly which document ids it will return. It sends a second, much smaller request to the shards owning those ids, asking for the stored fields in `fl`, plus per-document work such as highlighting. The results are assembled in the merged order and returned. The two-phase design exists to avoid transferring, serialising and highlighting documents that will lose the merge. On a wide document with large stored fields and highlighting across many shards, that saving is substantial. ## Single-pass When documents are small and `rows` is modest, the second round trip can cost more than the wasted transfer it prevents. `distrib.singlePass=true` tells the aggregator to request the full field list in phase one and skip phase two, trading extra bytes for one less network round trip. It is a per-request tuning knob, not a default. ## Distributed scoring BM25 and TF-IDF scores depend on collection-wide statistics: document frequency of a term, average field length, total document count. In a distributed query, each shard by default computes those statistics **from its own local index**. If term distributions differ across shards — which is common with routing prefixes, time-based data or simply small shards — the same document text can score differently depending on which shard it lives on, and the merged ranking is subtly wrong. Solr's answer is a configurable stats cache (`<statsCache>` in `solrconfig.xml`) with implementations that gather global term statistics before scoring, at the cost of an extra round trip per query. Most clusters live with local statistics because shards are large and reasonably uniform; the problem shows up on small or deliberately skewed shards, and it is worth naming as a known effect. ## Deep paging Because every shard must return `start + rows` results, `start=100000&rows=10` makes every shard build and ship a 100,010-entry list for the aggregator to merge and throw away. Cost grows linearly with `start` and multiplies by shard count. The fix is `cursorMark`, which pages by the last document's sort values instead of by an offset, so each request stays the same size no matter how deep you have gone. It requires a deterministic total ordering — a sort that includes the uniqueKey as a tiebreaker. ## Other traffic in the same request Faceting adds its own exchange: a shard can only report its local top facet values, so counts for a value that is popular globally but outside a shard's local top-N may be undercounted. Solr can issue a refinement round to ask shards specifically about candidate values, which costs another round trip but fixes the counts. ## Debugging `debug=track` shows the per-phase, per-shard timings, which is the fastest way to see whether a slow distributed query is slow everywhere or is waiting on one lagging shard. `distrib=false` against a specific replica isolates one shard's behaviour. And remember that a distributed query is as slow as its slowest shard — one oversized or hot shard sets the latency for the whole collection.

  • Why must each shard return start + rows results in the first phase rather than just rows?
    Because the aggregator has no idea how the global winners are distributed. All of the requested window could sit on a single shard, so every shard must offer enough candidates to cover that case. This is exactly why deep paging is expensive in a distributed search: with start=100000, each shard builds and ships a 100,010-entry candidate list that the aggregator merges and mostly discards.
  • When is distrib.singlePass=true a good idea, and when is it a bad one?
    Good when rows is small and the requested fields are few and short — you save a network round trip and the extra bytes are negligible. Bad when documents are large, the field list is wide, or highlighting is involved: phase one then transfers full documents from every shard, most of which lose the merge, so you pay far more bandwidth and serialisation than the round trip you saved.
  • Why can the same document score differently depending on which shard it lives on?
    Because scoring uses collection statistics — document frequency, average field length, document count — and each shard computes them from its own local index by default. Uneven term distribution across shards therefore produces inconsistent scores and a subtly wrong merged ranking. Solr can be configured with a stats cache implementation that gathers global statistics first, at the cost of an extra round trip per query.
  • How do you tell whether a slow distributed query is slow everywhere or waiting on one shard?
    Add debug=track, which reports the per-phase, per-shard timings for the request. If one shard's phase-one time dwarfs the rest, you have a hot or oversized shard rather than a query-shape problem — confirm by hitting that shard's replica directly with distrib=false. A distributed query's latency is set by its slowest shard, so one bad shard degrades the whole collection.

saying these in an interview costs you the question

  • Says every shard returns only rows documents in phase one
  • Thinks stored fields are fetched from all shards
  • Believes distributed scoring uses global statistics by default
  • Uses large start values for deep paging in a sharded collection
  • Assumes the aggregator is a fixed dedicated node

context