skip to content

Why does an Elasticsearch cluster with thousands of tiny shards perform worse than one with fewer large shards?

level: middleimportance: must knowfreq 70%

answer

  1. Each shard is a whole Lucene index
  2. Fixed overhead does not shrink with data
  3. Every shard adds a task to every search
  4. Heap and cluster state both pay per shard
  5. Guidance is a range, not 'fewer is better'

basics

~20 s

Every shard is a separate Lucene index with fixed overhead in heap, file handles, merge work and cluster-state metadata, and every search becomes a task per shard that must be scheduled and merged. Those per-shard costs dominate once shards are small, which is why guidance targets tens of gigabytes per shard.

solid answer

~50 s

A shard is not a free subdivision — it is a complete Lucene index. Each one carries its own segment metadata, open file handles, doc-values readers and merge threads, and each one adds entries to the cluster state that the master must publish to every node on change. On the search side, a query fans out to every shard it might match: each shard search is a separate task on the search thread pool, and the coordinating node has to merge all of those partial results. With thousands of 200 MB shards you pay all of that overhead to search very little data per task, so latency is dominated by scheduling and merging rather than by real work. Elastic's guidance is to keep most shards in the tens of gigabytes, stay well under roughly 20 shards per gigabyte of JVM heap on a node, and note that `cluster.max_shards_per_node` defaults to 1000 open shards per data node as a hard stop.

code

bash · 5 lines
bash
# shard size distribution, smallest first
GET _cat/shards?v&h=index,shard,prirep,docs,store&s=store

# total shards vs heap for the shards-per-GB check
GET _cluster/stats?filter_path=indices.shards.total,nodes.jvm.mem

go deeper

for a junior

Recall the shape of the guidance: a shard should hold tens of gigabytes, and creating hundreds of tiny shards makes a cluster slower, not faster.

for a middle

Explain the concrete costs — per-shard heap, file handles, merge threads, cluster-state entries and one search task per shard — and why they do not shrink as the shard shrinks.

for a senior

Diagnose it from the APIs: shard-size distribution, shards per GB of heap, search thread-pool rejections and slow master operations, then propose consolidation rather than per-index fiddling.

for a principal

Own the budget across the fleet: define shard-count and shard-size standards in index templates, keep total shards within heap-derived limits as data sources multiply, and treat unchecked index-per-source-per-day growth as an architectural defect.

## What a shard costs before it holds any data An Elasticsearch shard is a full Lucene index. That means it independently owns: - a set of **segments**, each with its own terms dictionary, postings, stored fields and doc values, and per-segment structures that must be opened and tracked; - **file handles** — a segment is several files, and a node with thousands of shards can run into operating-system file-descriptor limits; - **heap** for segment metadata, cached readers and per-shard bookkeeping; - **background merge work**, since every shard runs its own merge policy and merge threads; - **cluster-state entries**: the index metadata, mappings and routing-table entries for each shard copy. None of that scales down with the amount of data in the shard. Splitting one 50 GB shard into 250 shards of 200 MB does not reduce the total data — it multiplies the fixed overhead by 250. ## Cluster state and the master The cluster state holds index metadata, mappings and the routing table describing where every shard copy lives. It is kept in heap on every node and re-published by the elected master whenever it changes. A cluster with tens of thousands of shards, especially spread over many indices each with its own mapping, ends up with a large cluster state; publication takes longer, master operations (index creation, allocation decisions, node join/leave) get slower, and a master failover or full-cluster restart becomes noticeably painful. Oversharding is therefore not only a query-latency problem — it degrades the control plane. ## The search fan-out A search against an index is executed per shard. The coordinating node sends a shard-level request to one copy of each shard, each of those runs as a task on the target node's `search` thread pool (sized from the node's processor count, with a bounded queue), and the coordinating node then merges the per-shard results — reducing hits, and combining aggregation results. Two things go wrong when shard count is high: 1. **Scheduling dominates useful work.** A shard search over 200 MB may take a millisecond of real work but still needs a thread, a network round trip and a slot in the reduce. Query latency approaches the cost of coordinating hundreds of near-empty tasks. 2. **Thread-pool pressure.** Concurrent searches each multiply by shard count. Queues fill, and once full, Elasticsearch rejects shard requests — surfacing as search rejections and partial failures under load rather than graceful slowdown. Elasticsearch mitigates some of this: a pre-filter phase can skip shards that cannot match (particularly effective for time-range queries on time-based indices), and the `max_concurrent_shard_requests` request parameter bounds how many shard requests one search issues at a time. Those help; they do not make oversharding free. ## The rules of thumb, and where they come from - **Shard size in the tens of gigabytes.** Elastic's sizing guidance is roughly 10–50 GB per shard for most search and log workloads. Below that, fixed overhead dominates; far above it, recovery, rebalancing and snapshot restore get slow and coarse-grained because a shard is the unit that gets moved. - **Shards per node bounded by heap.** Aim to stay under about 20 shards per gigabyte of JVM heap on a data node. A node with 30 GB of heap should therefore stay well under about 600 shards. - **`cluster.max_shards_per_node`** defaults to 1000 open shards per data node. It is a cluster-wide budget — the limit multiplied by the number of data nodes — and exceeding it makes index creation fail with a validation exception rather than degrading silently. It is a safety net, not a target. All of these budgets are consumed by *shard copies*, not primaries: an index with 10 primaries and 1 replica costs 20 shards. ## The opposite failure Undersharding has its own problems, which is why the answer is a range and not "fewer is always better". A single enormous shard cannot be split across nodes, so it caps the parallelism of a query and the disk it can occupy; recovery of a 500 GB shard after a node failure moves a great deal of data; and rebalancing can no longer even out the cluster because the unit of movement is too coarse. ## Diagnosing it `GET _cat/shards?v&s=store` shows shard sizes and where they are; `GET _cat/indices?v` gives per-index totals; `GET _cluster/stats` reports total shard counts and heap. The classic oversharded signature is hundreds or thousands of indices, most of them a few hundred megabytes, one index per day per small data source. The fix is almost never to resize the existing indices one by one — it is to consolidate: fewer primaries in the template, roll over on primary shard size instead of on the calendar, and shrink or delete what already exists.

  • What does cluster.max_shards_per_node actually limit, and what happens when you hit it?
    It caps open shards — primaries plus replicas — per data node, defaulting to 1000, and is enforced cluster-wide as the limit multiplied by the number of data nodes. When the budget is exhausted, creating a new index or opening a closed one fails with a validation exception. It is a guard rail against runaway shard growth, not a performance target to aim at.
  • How would you spot an oversharded cluster from the APIs alone?
    `GET _cat/shards?v&s=store` sorted by size reveals a long tail of tiny shards; `GET _cat/indices?v` shows many small indices, typically one per day per source; `GET _cluster/stats` gives total shard count against total heap so you can check the shards-per-GB-heap ratio. Slow master operations and search thread-pool rejections corroborate it.
  • Can a shard be too large as well as too small?
    Yes. A shard is the unit of allocation, recovery, rebalancing and snapshot restore, so a very large shard means slow recovery after a node loss, coarse rebalancing that cannot even out disk usage, and a hard cap on how far the index can spread. That is why guidance is a band in the tens of gigabytes rather than a lower bound only.

Ten thousand one-page files in ten thousand folders take far longer to search than a hundred thick folders, even though the total number of pages is identical — the cost is in opening and closing folders, not in reading.

saying these in an interview costs you the question

  • Says more shards always means more parallelism and speed
  • Treats shards as free logical partitions with no fixed cost
  • Ignores replicas when counting shards against per-node budgets
  • Believes cluster.max_shards_per_node is a recommended target
  • Proposes one index per day per small data source without volume checks

context