Chapter 13 · Shard Sizing, Routing, Allocation, Awareness, and Cluster Topology

Document Routing, Custom Routing, Hotspots, Skew, and Co-Locating Queries

Use routing only when the query key and data distribution justify it, then prove that reduced fan-out does not create hot shards or correctness traps.

Intermediate115–145 minutesShard topology & failure-domain labElasticsearch 9.5.3 · OpenSearch 3.8.0Last reviewed: September 2026

Learning outcomes

AtlasMart adds enterprise tenants. Most requests filter by tenant_id, and the team considers custom routing so tenant searches touch fewer shards. Routing is powerful because it changes document placement, not because it is a query hint. The same choice can reduce fan-out for one workload while creating hot shards, duplicate IDs under different routing values, or missing documents when callers forget the routing key.

01

Explain default document routing and how a custom routing value selects a primary shard.

02

Prove that routed searches touch fewer shards and quantify the tradeoff.

03

Detect skew and hot routing keys before production saturation.

04

Protect get/update/delete correctness when custom routing is required.

05

Decide when routing is inferior to ordinary filtering, separate indices or a different partitioning strategy.

Chapter baseline reviewed 11 September 2026

Examples target self-managed Elasticsearch 9.5.3 / Kibana 9.5.3 and OpenSearch 3.8.0 / OpenSearch Dashboards 3.8.0. The established AtlasMart endpoints remain https://localhost:9200 for Elasticsearch using its copied CA and https://localhost:9201 for the disposable OpenSearch demo certificate. Use each distribution's bundled JVM for this lab unless its current support matrix says otherwise. The initial Chapter 01 single-node containers are useful for API inspection but cannot demonstrate replica placement or zone resilience, so this chapter uses deterministic traces and an optional disposable multi-node topology. No moving latest tags and no production allocation changes are used.

Execution and safety note

Run mutating, destructive, security, lifecycle, snapshot, failure-injection, and load-test commands only in the disposable AtlasMart lab or an equivalently isolated environment. Verify the target cluster, index, tenant, credentials, and rollback path before execution; treat shown output as an expected invariant unless the lesson explicitly labels it as captured evidence.

1. Routing chooses placement; filtering chooses matches

By default, Elasticsearch and OpenSearch derive a routing value from the document ID and hash it into the index’s routing space. With custom routing, the application supplies a stable value such as tenant_id. Documents with the same routing key are intentionally co-located into a subset—often one—of the primaries.

A filter runs after the request has reached candidate shards. A routing parameter can reduce which shards receive the request at all. That can reduce fan-out, but it also means every read/write path must preserve the routing contract.

Dev Tools · create a disposable routed index
PUT atlasmart-routing-lab
{
  "settings":{"number_of_shards":4,"number_of_replicas":0},
  "mappings":{
    "_routing":{"required":true},
    "properties":{
      "tenant_id":{"type":"keyword"},
      "sku":{"type":"keyword"},
      "name":{"type":"text"},
      "price":{"type":"scaled_float","scaling_factor":100}
    }
  }
}

2. Make placement and fan-out observable

Index with explicit routing
PUT atlasmart-routing-lab/_doc/P-2001?routing=tenant-a
{"tenant_id":"tenant-a","sku":"P-2001","name":"USB-C Dock","price":89.00}

PUT atlasmart-routing-lab/_doc/P-2002?routing=tenant-b
{"tenant_id":"tenant-b","sku":"P-2002","name":"4K Monitor","price":329.00}

GET atlasmart-routing-lab/_search_shards?routing=tenant-a
GET atlasmart-routing-lab/_search?routing=tenant-a
{"query":{"term":{"tenant_id":"tenant-a"}}}

Compare _search_shards with and without routing=tenant-a. The routed form should identify a smaller shard set for this 4-primary fixture. The exact shard number is an implementation result of hashing; do not hard-code it into application logic.

Routing is not an authorization boundary. A caller with direct index access and no tenant filter can still search other shards. Tenant isolation belongs in authorization plus tested query construction, not in shard placement alone.

3. The missing-routing and duplicate-ID boundary cases

When _routing.required=true, an index request without routing is rejected, which is useful because it turns a silent placement mistake into an explicit failure. The application must carry the same routing value for GET, update and delete operations.

Custom routing also changes the uniqueness scope of an ID in practice: the same _id combined with different routing values can address different shard locations. Treat (routing, _id) as the application identity contract and prevent accidental duplicates upstream.

Controlled failure · omitted routing
PUT atlasmart-routing-lab/_doc/P-2003
{"tenant_id":"tenant-a","sku":"P-2003","name":"Laptop Stand","price":39.00}

# Expected shape: a routing-missing error because the mapping requires _routing.
# Do not memorize the full message text across versions.
Correctness rule

If a write used custom routing, every point read/update/delete must know the routing value. Do not “repair” a missing result by issuing an unrestricted search and then mutating whatever document happens to match.

4. Hotspots are a distribution problem, not a shard-count problem

If one tenant contributes 45% of writes, routing by tenant can send that traffic to one shard even when the index has many primaries. Adding shards does not guarantee that the heavy tenant spreads because the routing contract intentionally co-locates it. This is the classic hot-key problem.

Build a routing-key histogram before enabling custom routing, then repeat it over time. Observe per-shard indexing/search totals, CPU and queueing rather than only cluster-wide averages.

Python · deterministic routing-key skew report
from collections import Counter

keys = (["tenant-a"] * 4500 + ["tenant-b"] * 1500 +
        ["tenant-c"] * 1200 + [f"tenant-{i}" for i in range(1000, 3800)])
counts = Counter(keys)
total = sum(counts.values())
for key, n in counts.most_common(5):
    print(key, n, f"{100*n/total:.1f}%")

# Feed real request/indexing logs into the same shape.
# A dominant key is a design signal; no universal percentage threshold is safe.
Observe shard-level activity
GET atlasmart-routing-lab/_stats/indexing,search,store?level=shards
GET _cat/shards/atlasmart-routing-lab?v&h=index,shard,prirep,state,docs,store,node

5. Deliberately wrong approach: route every workload by tenant

Tenant routing is attractive when most queries are tenant-scoped, but it can hurt cross-tenant analytics because those requests still need all relevant shards. It can also concentrate a whale tenant and make one shard the limiting resource. Routing is therefore a workload-specific optimization with a data-placement cost.

Repair the design by considering alternatives: ordinary tenant filters when fan-out is acceptable; separate indices for a small number of truly isolated large tenants; time/data-stream partitioning for append-heavy workloads; or an application-level partition key that distributes a heavy tenant while preserving the dominant query path.

Pattern Good fit Watch for
Default routing Mixed/global queries; balanced IDs All-shard fan-out for broad index searches.
Tenant custom routing Mostly tenant-scoped queries; bounded tenant sizes Hot tenants; routing-key propagation.
Tenant + bucket key One tenant is too large for one shard Queries may need multiple bucket routing values.
Separate index Strong lifecycle/security/topology isolation Many indices/tiny-shard overhead if overused.

6. AtlasMart lab: compare uniform and skewed routing

Create two synthetic distributions with the same document count: one uniform across many tenants and one dominated by a few tenants. Index them into disposable 4-primary indices using a stable routing field. Compare shard document counts and per-shard indexing totals. Then run the same tenant query with and without routing and record touched shard count plus latency distributions under identical concurrency.

The lab passes only if you can explain both the fan-out reduction and the worst-shard imbalance. A lower average query latency does not justify a design that violates write/recovery headroom on one shard.

  • Verify every routed write has an explicit routing key.
  • Verify point reads fail or miss as expected when the routing key is omitted/wrong.
  • Record max/min shard document and byte ratios for both distributions.
  • Measure tenant-scoped and global query latency separately.
  • Delete only the disposable atlasmart-routing-lab* indices during cleanup.

Production judgment

Custom routing is a placement contract that should be versioned with the application API. Log routing-key distribution, reject missing keys early, and include routing behavior in migration/reindex tests. If a future query pattern becomes global, the original routing optimization may no longer be a win.

Check your understanding

  1. What changes when custom routing is supplied?
  2. Why can a routed query be faster?
  3. What is the hot-key risk?
  4. Is routing a tenant-security boundary?
  5. Why require routing in the mapping?
Review the answers

1. The routing value participates in shard selection, changing where the document is stored and which shards a routed request targets.

2. It may execute on fewer shards, reducing per-shard work and coordination when the routing key is correct.

3. A high-volume routing value can concentrate writes, searches and storage on one shard.

4. No. Authorization and query constraints must enforce tenant isolation independently.

5. It converts a missing routing key from a silent placement bug into a rejected write.

Summary and next step

You can now treat routing as a measurable data-placement choice rather than a magic speed flag. Lesson 3 moves down a layer to the allocator: why a shard can or cannot live on a node, how disk watermarks interact with awareness, and what a zone failure should look like.

Authoritative references

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.