skip to content

How does SolrCloud's compositeId router decide which shard a document lands on?

level: middleimportance: should knowfreq 50%

answer

  1. Something is hashed into a fixed range
  2. Shards own slices of that range
  3. The exclamation mark in an id matters
  4. High bits versus low bits
  5. A parameter narrows the fan-out at query time

basics

~20 s

It hashes the document's uniqueKey into a 32-bit value and sends the document to whichever shard owns that value's hash range. If the key contains a routing prefix such as tenant!docId, the high bits come from the prefix, so every document sharing that prefix lands on the same shard.

solid answer

~60 s

With the default `compositeId` router, each shard of a collection owns a contiguous slice of a 32-bit hash space — you can see the slices as `range` entries in the collection's `state.json`. Solr hashes the document's uniqueKey and routes the document to the shard whose range contains the resulting value. Because the ranges are fixed at creation, `numShards` is required when you create a compositeId collection. The "composite" part is the routing prefix. If the id is written `acme!inv-42`, Solr takes the high bits of the hash from `acme` and the low bits from `inv-42`. All of tenant `acme`'s documents therefore fall in one narrow, contiguous hash band and land on a **single shard**. At query time `_route_=acme!` restricts the request to that shard, turning a fan-out over every shard into a single-shard query. The cost is skew: a large tenant becomes a hot shard. The `key/bits!id` form lets you take fewer bits from the prefix so a big tenant is deliberately spread across a wider slice of the hash space.

code

json · 9 lines
json
{
  "products": {
    "router": { "name": "compositeId" },
    "shards": {
      "shard1": { "range": "80000000-ffffffff", "state": "active" },
      "shard2": { "range": "0-7fffffff", "state": "active" }
    }
  }
}

go deeper

for a junior

Know that Solr hashes the uniqueKey and each shard owns a hash range, and that an id like tenant!id keeps a tenant's documents together on one shard.

for a middle

Explain the high-bits/low-bits construction, why numShards is fixed at creation, and what route does to the fan-out on a query.

for a senior

Diagnose skew in production: spot the hot shard, know the bit-mask form and its reindex cost, and recognise when SPLITSHARD will not help because a prefix pins the high bits.

for a principal

Decide the multi-tenant layout before data exists — shared collection with prefixes, bit-masked large tenants, or dedicated collections behind aliases — and account for the reindex you are committing to if the choice turns out wrong.

## The hash space When a collection is created with the default `compositeId` router, Solr divides a 32-bit hash space into `numShards` contiguous ranges and writes them into `/collections/<name>/state.json`. A two-shard collection looks like `shard1: 80000000-ffffffff`, `shard2: 0-7fffffff`. Indexing a document means hashing its uniqueKey to a 32-bit value and delivering it to the shard whose range contains that value; the receiving node forwards it to that shard's leader. Querying without a routing hint means fanning out to one replica of every shard, because any shard could hold a match. This is why `numShards` is mandatory for a compositeId collection and why you cannot simply "add a shard" to one: the ranges must partition the space exactly. To add capacity you SPLITSHARD, which halves an existing shard's range into two sub-shards and retires the parent. ## The composite part The router is called *composite* because the id may be composed of a routing prefix and a document id, separated by `!`: ``` acme!inv-42 ``` Instead of hashing the whole string as one unit, Solr hashes the prefix and the remainder separately and builds the final 32-bit value from the **high bits of the prefix hash** and the **low bits of the document-id hash**. Every document whose prefix is `acme` therefore shares the same high bits and lands inside one narrow contiguous band of the hash space — which, unless that band happens to straddle a boundary, means one shard. The payoff is at query time. `_route_=acme!` tells the aggregator which shard owns that prefix, so the request goes to exactly one shard instead of all of them. On a 20-shard collection that removes 19 shard requests, 19 sets of top-N merging, and the tail latency of the slowest shard. The same parameter works for deletes and for atomic updates. ## Bit masking for uneven tenants Co-location is only a win while the co-located set is a reasonable size. Give one tenant 200 million documents in a cluster whose average shard holds 20 million, and that shard becomes a hot spot no amount of replication fixes — replicas share the query load but every one of them still stores the whole oversized shard. For that case compositeId supports a bit count: `acme/4!inv-42` says "take only 4 bits from the prefix hash". Fewer prefix bits mean the tenant's documents occupy a *wider* slice of the hash space and therefore spread over more shards, while queries with `_route_=acme/4!` still restrict the fan-out to just that slice rather than the whole collection. It is a dial between perfect co-location and even distribution — but it must be chosen up front, because changing the bit count changes where documents hash to and requires a reindex of that tenant. CompositeId also supports a second level (`app!user!docId`), which lets you route at either granularity; the high half of the hash is shared between the two prefixes. ## The implicit router The alternative is `router.name=implicit`: Solr does no hashing at all, and the document goes to the shard you name — either through the `_route_` parameter on the request or through a field on the document named by `router.field`. You supply the shard names with the `shards` parameter at creation, and you can add more later with CREATESHARD (which is only allowed for implicit collections). This suits time-based or region-based partitioning where you want explicit control over which slice a document lands in, at the price of owning balance yourself. ## Failure modes to name in an interview - **Routing prefix baked into the id.** The prefix is part of the uniqueKey. `acme!inv-42` and `inv-42` are different documents, and moving a document between tenants means delete-and-reindex. - **Forgetting `_route_` on reads.** Indexing with prefixes but querying without them gives you all the skew and none of the speed-up. - **`_route_` on a query that must see everything.** Restricting the fan-out means results from other shards are silently missing, not merely deprioritised — a correctness bug, not a performance one. - **Skew after growth.** A tenant that was small at design time becomes the hot shard later; monitor per-shard document counts and index size, not just cluster totals. - **Assuming SPLITSHARD fixes prefix skew.** Splitting a shard whose documents all share one prefix's high bits may put nearly everything into one sub-shard, because the prefix pins the high bits that the split boundary uses.

  • Why is numShards required when creating a compositeId collection, but not for an implicit one?
    Because compositeId partitions a 32-bit hash space into contiguous ranges that must cover it exactly; Solr needs the shard count to compute those ranges at creation and record them in state.json. An implicit collection does no hashing — you name the shards with the `shards` parameter and route documents explicitly — so there is no range arithmetic to do, and you may add shards later with CREATESHARD.
  • What goes wrong if you index with routing prefixes but query without _route_?
    Nothing breaks functionally — the fan-out to every shard still finds the documents — but you pay all the cost of co-location and get none of the benefit. Your tenants are unevenly distributed across shards, so some shards are large and hot, and every query still contacts all of them and waits for the slowest. Prefix routing is only a win when reads carry the matching `_route_`.
  • A single tenant grows to ten times the size of any other. What are your options?
    Reindex that tenant with a bit-masked prefix such as `bigtenant/4!id`, so its documents spread over a wider hash band while `_route_=bigtenant/4!` still limits fan-out. Alternatives are giving that tenant its own collection behind an alias, or SPLITSHARD on the affected shard — though splitting helps little when one prefix pins the high bits that define the split boundary.
  • Does _route_ on a query change scoring or just which shards are contacted?
    Only which shards are contacted. The matched documents are scored exactly as they would be on that shard anyway. It does have a side effect on distributed scoring, though: with the default local term statistics, restricting to one shard removes the cross-shard IDF inconsistency, because all scores now come from one set of local statistics.

saying these in an interview costs you the question

  • Thinks the routing prefix is stored separately from the id
  • Believes _route_ merely prioritises shards rather than restricting them
  • Says you can add a shard to a compositeId collection with CREATESHARD
  • Assumes co-location is free regardless of tenant size
  • Expects hashing to balance a workload with one dominant tenant

context