skip to content

How do you choose the primary shard count for a new Elasticsearch index whose future size you don't know?

level: principalimportance: should knowfreq 45%

answer

  1. Measure one shard before choosing many
  2. Volume divided by a measured budget
  3. Two bounds: node spread below, overhead above
  4. The best answer avoids a permanent bet
  5. Write down the assumptions with the number

basics

~20 s

Benchmark one shard loaded with representative data and real queries to find the size where latency degrades, then divide the projected data volume by that size. Where growth is genuinely unknown, design so the number never has to be right forever: roll over to new indices sized on shard size.

solid answer

~50 s

Start empirically rather than from a formula. Build a single-shard index, load representative documents, run the real query mix, and grow it until latency crosses your service objective — that gives you a defensible shard size for *this* workload, which usually lands somewhere in the tens of gigabytes. Divide your projected data volume by it and round with headroom for growth and for replicas. Then apply the constraints: primary count caps how many nodes the index can spread over, so it must be at least the number of data nodes you intend to use, while an inflated count multiplies heap, cluster-state and per-search fan-out costs. The strategic answer is to avoid a permanent bet: for append-only data, use rollover so each new index gets a fresh, correctly-sized decision, and keep reindex, `_split` and `_shrink` as the escape hatches for the indices where the number really is fixed.

code

bash · 6 lines
bash
# single-shard sizing rig: load, query, watch p99 as store grows
PUT /sizing-probe
{ "settings": { "number_of_shards": 1, "number_of_replicas": 0 } }

GET _cat/shards/sizing-probe?v&h=shard,docs,store
GET /sizing-probe/_stats/search?filter_path=**.query_time_in_millis

go deeper

for a junior

Know that the shard count is chosen at index creation and cannot be changed, and that the aim is shards in the tens of gigabytes rather than very many tiny ones.

for a middle

Explain the arithmetic: measured per-shard capacity divided into projected volume, and the two bounds — enough shards to spread across nodes, few enough to keep per-shard overhead down.

for a senior

Show the method: benchmark a single shard with representative data and the real query mix, project growth with retention, and preserve an escape hatch such as an alias plus reindex.

for a principal

Own the strategy: make rollover and shared index templates the default so no single sizing decision is permanent, define the fleet-wide shard budget against heap, and document the assumptions behind every standard so the number is auditable.

## Why this is a judgment question `index.number_of_shards` is final. Everything about the decision follows from that: you are picking a number that will govern an index for its entire life, usually before you have seen production data or production query patterns. There is no formula that produces the right answer, so the real skill is designing so that being wrong is cheap. ## Step one: measure, do not guess The only number that matters is how much data a single shard of *your* documents, with *your* mappings, can hold while still meeting *your* latency objective. Get it like this: 1. Create a one-primary, zero-replica index with the production mapping and analysis chain. 2. Load a representative sample — same field cardinality, same document size, same proportion of nested documents or vectors, because all of those change the answer dramatically. 3. Replay a realistic query mix, including the expensive aggregations and the deepest paginations, at realistic concurrency. 4. Grow the shard and watch p99 latency, heap and merge activity. The size where latency crosses your objective is your shard-size budget. A log workload with keyword filters and date ranges will tolerate a very different shard size from a product catalogue with heavy aggregations and highlighting, or an index carrying dense vectors. This is why generic guidance is a broad band in the tens of gigabytes rather than a single number: it is the range most workloads land in, not a substitute for measuring yours. ## Step two: project the volume and divide Estimate data volume at the retention horizon you actually keep — daily ingest multiplied by retention, plus the index overhead over raw document size, which you now know from the benchmark. Divide by your measured shard budget, then round up for growth headroom. Remember to count replicas when checking disk and node capacity: each replica is another full copy. ## Step three: apply the structural constraints Two hard bounds frame the number: - **Lower bound: node spread.** An index cannot occupy more data nodes than it has shard copies. If you intend to spread an index across eight data nodes, one primary with one replica cannot do it. Shard count is also what allows a single query's work to be spread across machines. - **Upper bound: overhead.** Every shard copy costs heap, file handles, merge threads, cluster-state entries, and one task per search. Guidance to stay under roughly 20 shards per gigabyte of JVM heap on a node, and the `cluster.max_shards_per_node` default of 1000 open shards per data node, are the guard rails. An inflated count buys nothing and degrades every query. Between those bounds, prefer the smaller number that still meets the spread requirement. Over-sharding "just in case" is the more common and more expensive mistake, because its cost is paid on every single search forever, whereas under-sharding is fixable with a split or a reindex. ## Step four: design so the bet is not permanent This is what separates a principal-level answer: - **Append-only or time-series data** (logs, metrics, events, orders) should not live in one growing index at all. Roll over to a new index on a primary-shard-size condition, keep the shard count in a shared template, and let retention delete the old ones. The shard count is then revisited implicitly on every roll, and a wrong choice ages out instead of becoming permanent. - **Mutable corpora** (a product catalogue, a user directory) do live in one index, so sizing matters more — but they are also usually small enough that one or a few shards is right, and a reindex behind an alias is a routine, rehearsed operation. Always front such an index with an alias from day one, so the swap costs nothing when the day comes. - **Set `index.number_of_routing_shards` deliberately** if you expect an index to need splitting later; it is final too, and it caps the split factors available to you. ## Common bad heuristics to challenge - *"One shard per data node."* It guarantees spread but ignores data volume entirely, and it turns node additions into a shard-count problem. - *"Pick the maximum you might ever need — extra shards are free."* They are not; the fixed per-shard overhead is paid whether the shard holds 200 MB or 50 GB. - *"We will just split it later."* Split requires a write block and a compatible routing-shard configuration; it is a planned maintenance operation, not a live resize. - *"Match the shard count to CPU cores."* Search parallelism is bounded by shard count, but so is nothing else about capacity; the data volume is the primary driver. ## How to present the decision Write it down: the measured shard budget, the volume projection and its assumptions, the resulting count, the node-spread requirement it satisfies, and the escape hatch you have preserved (alias plus reindex, rollover, or a routing-shard configuration that allows splitting). A shard count with that reasoning attached is defensible; a number with no reasoning is the thing that gets copied into the next team's template and causes the next oversharded cluster.

  • What lower bound does node spread place on the shard count?
    An index cannot occupy more data nodes than it has shard copies, so if you want it spread over eight nodes you need at least eight copies — for example four primaries with one replica each. Shard count is also the ceiling on how many machines can share the work of one query against that index, so a single-shard index is pinned to one node's capacity.
  • When is picking a permanent shard count unavoidable, and how do you de-risk it?
    For a mutable corpus that must live in one index — a catalogue, a user directory — rollover does not apply. De-risk it by always addressing the index through an alias so a reindex-and-swap is a routine operation, setting `index.number_of_routing_shards` deliberately so a split stays available, and rehearsing the reindex before you need it.
  • Why is over-provisioning shards worse than under-provisioning?
    Over-provisioning costs heap, cluster-state size and one extra search task per shard on every query, forever, and the only fix is a shrink or reindex. Under-provisioning is usually caught by the shard-size metric before it hurts, and split, shrink and reindex all exist to correct it. Given a symmetric uncertainty, err smaller.

saying these in an interview costs you the question

  • Picks one shard per data node with no volume analysis
  • Over-provisions shards because they are assumed free
  • Assumes the count can be tuned later without a new index
  • Sizes on raw document bytes, ignoring index and replica overhead
  • Reuses another cluster's shard count without benchmarking

context