skip to content

How would you decide whether to use custom routing for a multi-tenant Elasticsearch index with thousands of tenants?

level: principalimportance: should knowfreq 32%

answer

  1. Compare fan-out cost against skew cost
  2. Today every tenant search touches every shard
  3. A shard is the smallest rebalance unit
  4. Whale tenants deserve their own index
  5. Spread a tenant over k shards instead of one

basics

~20 s

Weigh the fan-out saved against the skew created. Custom routing turns every tenant search into a single-shard search, which multiplies search throughput, but it concentrates each tenant on one shard, so an outsized tenant becomes a hot shard the cluster cannot rebalance away.

solid answer

~50 s

Start from the tenant size distribution, not from the feature. Without routing, every tenant's documents are spread over every shard, so a search for one tenant fans out to all of them and per-query cost scales with shard count rather than with the tenant's data. With `?routing=<tenant>`, a tenant's documents share one shard and a tenant-scoped search touches one shard — throughput rises with node count and tail latency drops, since you no longer wait on the slowest of N shards. The costs are real: the routing value must accompany every index, get, update, delete and search, so declare `_routing` required and hide the value behind per-tenant aliases with `routing`/`search_routing`; and a shard is the smallest unit the cluster can move, so one whale tenant produces a shard nothing can rebalance. The usual answer is a tiering strategy — dedicated indices for the few whales, a routed shared index for the long tail, and `index.routing_partition_size` to spread mid-size tenants across a fixed handful of shards.

code

json · 10 lines
json
{
  "actions": [
    { "add": {
        "index": "orders",
        "alias": "orders-tenant-7",
        "filter": { "term": { "tenant": "tenant-7" } },
        "routing": "tenant-7"
    } }
  ]
}

go deeper

for a junior

Know that adding a routing value on indexing and searching sends a tenant's documents to a single shard, and that the same value must be used consistently.

for a middle

Explain the mechanism and its contract: the routing value replaces the id as the hash input, so every get, update and delete must repeat it or the document appears missing.

for a senior

Diagnose and operate it: spot the hot shard produced by an outsized tenant, enforce the routing contract with a required _routing mapping and per-tenant aliases, and know the partition-size escape hatch.

for a principal

Own the strategy — measure the tenant distribution and query mix, design a tiered layout with aliases as the seam so tenants can be promoted between shared and dedicated indices, and state which failure you are choosing to accept.

## Frame the decision as a cost comparison Two costs are in tension, and the tenant size distribution decides which dominates. **Fan-out cost.** Without custom routing, hashing on `_id` spreads every tenant's documents across every shard. A search scoped to one tenant must still ask every shard, because any shard could hold a match. Each shard-level request consumes a slot in the node's search thread pool, and the coordinating node cannot answer until the slowest of them replies. Per-query cost is proportional to the shard count, independent of how small the tenant is. With thousands of small tenants querying constantly, the search tier saturates on coordination overhead rather than on real work. **Skew cost.** With custom routing, a tenant's documents all hash to one shard. Shard sizes now follow your tenant size distribution, which in every real multi-tenant product is long-tailed. The largest tenant defines the largest shard, and a shard is atomic: the cluster can move it between nodes but cannot split one tenant's data across two. No amount of rebalancing helps. ## Measure before choosing The inputs you should bring to this decision: - **Tenant size distribution** — median, p99 and maximum documents and bytes per tenant. A distribution where the largest tenant is 50× the median is a different problem from one where it is 5,000×. - **Query mix** — what fraction of searches are tenant-scoped versus cross-tenant? Routing only pays where searches carry the routing key. Global admin search, cross-tenant analytics and background jobs still fan out. - **Query rate per tenant** — skew in *traffic* is as damaging as skew in data. A small tenant with enormous query volume also creates a hot shard. - **Growth and churn** — tenants that grow 10× in a quarter turn a comfortable shard into a problem, and routing offers no incremental relief. - **Lifecycle requirements** — deleting a tenant from a routed shared index means a delete-by-query and the deleted documents linger until merges reclaim them; deleting a dedicated index is instantaneous. ## The operational contract routing imposes Once documents carry a custom routing value, that value is part of their identity for every future operation. A get, update or delete by id without it hashes the `_id`, lands on the wrong shard, and reports the document missing — silently. Two mechanisms make this safe: ``` "mappings": { "_routing": { "required": true } } ``` turns an omitted routing value into an error rather than a misrouted write. And a per-tenant alias can carry both a filter and the routing values, so application code targets `orders-tenant-7` and cannot forget either: ``` { "index": "orders", "alias": "orders-tenant-7", "filter": { "term": { "tenant": "tenant-7" } }, "routing": "tenant-7" } ``` The alias's `routing` sets both index-time and search-time routing; `index_routing` and `search_routing` set them separately when they must differ. ## The middle options This is rarely a binary choice. - **`index.routing_partition_size`.** Set to *k*, a routing value maps to *k* shards rather than one: the shard is chosen from the routing hash combined with the `_id` hash. A tenant's searches then touch *k* shards instead of all of them, while its data spreads over *k* shards instead of piling onto one. It is a static setting, must be smaller than the primary count, requires routing to be mandatory, and does not coexist with join fields. It is the right tool for mid-size tenants where one shard is too small and full fan-out is too expensive. - **Tiering by size.** Give the handful of whale tenants their own indices — full fan-out inside a dedicated index is cheap when that index has few shards, and lifecycle operations become index operations. Put the long tail in one routed shared index. Hide the split behind aliases so the application never knows which tier a tenant is in, and make promotion between tiers a reindex you can run on a schedule. - **No routing at all**, if searches are dominated by cross-tenant work or if the index is small enough that fan-out is a handful of shards. ## Second-order effects worth naming A tenant-scoped search on a routed index touches one shard, so its relevance scoring uses one consistent set of term statistics — an accidental improvement over the fan-out case. Conversely, a cross-tenant search over a routed index now faces a deliberately skewed document distribution, which is exactly the condition under which per-shard statistical skew distorts ranking. And the pre-filter that lets time-based searches skip shards does nothing for tenant fan-out, because every shard plausibly holds matching data. ## What interviewers listen for The answer that reads as principal-level starts by asking for the tenant size distribution and query mix instead of answering yes or no, names the hot-shard failure mode as unfixable by rebalancing rather than as a tuning problem, treats the routing value as a permanent contract with an enforcement mechanism, and proposes a tiered design with aliases as the seam — because the real decision is not "routing or not" but "how do we keep the option to change our mind per tenant".

  • What does index.routing_partition_size change, and what does it require?
    It makes a routing value map to a fixed number of shards instead of exactly one, by combining the routing hash with the document's own id hash. A tenant's searches then touch that many shards rather than all of them. It is static, must be less than the primary shard count, requires routing to be mandatory in the mapping, and is incompatible with join fields.
  • How do you stop application code from forgetting the routing value?
    Two layers. Declare `_routing` required in the mapping so a write without it fails loudly instead of misrouting. Then expose each tenant through an index alias carrying that tenant's `routing` (or `index_routing` and `search_routing`) plus a tenant filter, so application code targets the alias and the routing value is applied for it.
  • A routed shared index has one tenant occupying 40% of one shard's data. What are your options?
    Move that tenant out to its own index and repoint its alias — a reindex you can run online, after which the shared shard shrinks on merge. Rebalancing cannot help, because a shard is atomic and split multiplies all shards rather than relieving one. Longer term, add a size threshold that promotes tenants to dedicated indices automatically.
  • When is custom routing simply the wrong tool for a multi-tenant index?
    When most searches are cross-tenant, since those still fan out and gain nothing while inheriting the skew; when tenant sizes differ by orders of magnitude and the largest cannot fit comfortably in one shard; or when the index is small enough that fan-out already touches only a few shards and the operational contract buys no measurable latency.

saying these in an interview costs you the question

  • Proposes routing without asking about tenant size distribution
  • Thinks rebalancing will even out a hot routed shard
  • Forgets that gets and deletes need the routing value too
  • Suggests one index per tenant for thousands of tenants
  • Assumes routing helps cross-tenant searches as well

context