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.
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.
Explain default document routing and how a custom routing value selects a primary shard.
Prove that routed searches touch fewer shards and quantify the tradeoff.
Detect skew and hot routing keys before production saturation.
Protect get/update/delete correctness when custom routing is required.
Decide when routing is inferior to ordinary filtering, separate indices or a different partitioning strategy.
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.
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.
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
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.
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.
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.
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.
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
- What changes when custom routing is supplied?
- Why can a routed query be faster?
- What is the hot-key risk?
- Is routing a tenant-security boundary?
- 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
- Elastic shard allocation settings — Current allocation, rebalance, disk-watermark and recovery controls.
- Elastic shard allocation awareness — Awareness and forced-awareness behavior across failure domains.
- Elastic allocation explain API — Explain why a shard is or is not allocatable.
- Elastic index recovery API — Recovery stage and byte/file progress.
- OpenSearch cluster settings — Current routing/allocation, awareness and recovery settings.
- OpenSearch cluster tuning — Shard allocation awareness and forced awareness concepts.
- OpenSearch CAT shards — Shard state and placement inspection.
- OpenSearch cluster allocation explain — Allocation diagnostics and decider evidence.
- Elastic mapping routing field — Custom routing and required-routing behavior.
- OpenSearch search shard routing — Search routing parameters and distributed execution context.