Scaling Elasticsearch Clusters
Elasticsearch is distributed by design. Each index is split into shards spread across nodes, and replicas provide redundancy and extra read capacity. That lets a cluster grow from a single node on a laptop to hundreds of nodes holding petabytes of logs. But capacity doesn't come automatically: poor shard sizing, unbounded indices, slow bulk ingestion, and missing lifecycle policies are behind most Elasticsearch performance and stability problems.
This page covers the operational concepts that decide whether a cluster stays healthy: node roles, shard and replica design, time-series patterns with data streams and index lifecycle management (ILM), data tiers, indexing throughput, and how to keep search indices in sync with a primary database.
TL;DR
- An index has primary shards (fixed at creation, unless you split, shrink, or reindex) and replicas (changeable anytime).
- Aim for shards of roughly 10–50 GB (up to ~200M documents); avoid thousands of tiny shards.
- Separate node roles in larger clusters: dedicated master-eligible nodes (3), data nodes by tier, ingest and coordinating nodes.
- For logs, metrics, and events, use data streams plus ILM with rollover and hot → warm → cold → frozen → delete tiers.
- Index with the Bulk API, tune the refresh interval, and use aliases for zero-downtime reindexing.
- Monitor cluster health (green, yellow, red), heap, disk watermarks, and search and indexing latency.
Quick Example
An ILM policy and data stream for application logs:
Bulk indexing from Python:
Core Concepts
Nodes and Roles
Small clusters combine roles; larger clusters separate them so heavy searches don't destabilize master nodes. Master elections use a quorum, as in distributed consensus.
Shards and Replicas
- A primary shard is a self-contained Lucene index. The number of primaries is set at creation. Changing it requires the Split/Shrink APIs or reindexing.
- Replicas copy primaries onto other nodes for high availability and read throughput. A replica never sits on the same node as its primary.
- Search fans out to one copy of each shard and merges results. Indexing goes to the primary, then to its replicas.
Shard Sizing
Each shard has overhead (heap, file handles, cluster state), so too many small shards (oversharding) is the most common scaling mistake. Guidelines:
- Target roughly 10–50 GB per shard, and under about 200M documents.
- Keep the total shard count per node reasonable (Elastic recommends staying well under 1,000 non-frozen shards per node, and fewer is better).
- Use rollover by size and age for time-series data, so shard size stays consistent.
- Small, static indices (a product catalog of a few GB) usually need just one primary shard.
See database sharding.
Cluster Health
- Green: all primaries and replicas are assigned.
- Yellow: all primaries are assigned, some replicas aren't (normal on single-node clusters, and a risk elsewhere).
- Red: at least one primary is unassigned, so some data is unavailable.
GET _cluster/health, GET _cat/shards, and GET _cluster/allocation/explain diagnose issues. Disk watermarks (85%, 90%, 95% by default) block shard allocation, then writes, as disks fill, which is a frequent source of red and read-only indices.
Data Streams, ILM, and Data Tiers
For append-only time-series data (logs, metrics, traces, events):
- Data streams write to a single name while rolling over to new backing indices behind it.
- ILM automates rollover, then moves indices through tiers:
- Hot: fast SSD nodes for recent data, with indexing and frequent queries.
- Warm: cheaper nodes; shrink and force-merge older indices.
- Cold / frozen: searchable snapshots on object storage (S3, GCS) with a tiny local footprint.
- Delete: enforce data retention.
This keeps shards well-sized and costs proportional to data value. See log aggregation.
Indexing Throughput
- Use the Bulk API with batches of about 5–15 MB, and tune concurrency by measuring.
- Increase
refresh_interval(for example 30s, or-1during large backfills) to reduce segment churn, since near-real-time visibility costs throughput. - Temporarily set replicas to 0 during initial bulk loads, then restore them.
- Use
createoperations and auto-generated IDs for append-only data (it's faster), and handle per-document errors in bulk responses.
Keeping Elasticsearch in Sync
Elasticsearch is usually a secondary index, not the system of record. Sync patterns:
- Dual writes from the app (simple, but risk inconsistency when one write fails).
- Change data capture from the database (Debezium → Kafka → Elasticsearch sink), which is reliable and ordered. See Kafka Connect and change data capture.
- Periodic reindexing from the source for full rebuilds, combined with aliases for zero-downtime swaps.
Best Practices
Take Snapshots
Register a snapshot repository (S3, GCS, Azure) and schedule snapshot lifecycle management (SLM). Replicas aren't backups: they replicate deletions and corruption too.
Right-Size Heap and Memory
Give the JVM heap no more than about 50% of RAM (and under ~31 GB for compressed object pointers), leaving the rest for the OS file cache, which Lucene relies on heavily.
Use Managed Services or Operators
Elastic Cloud, Amazon OpenSearch Service, or ECK (Elastic Cloud on Kubernetes) handle upgrades, snapshots, and scaling. Self-managing large clusters requires real expertise.
Monitor the Right Signals
Track cluster health, JVM heap and GC, disk usage against watermarks, search and indexing latency and rejections (thread pool queues), and shard counts. Alert before disks hit watermarks. See monitoring.
Common Mistakes
Oversharding
Daily indices with 5 primaries and 1 replica for a service producing 200 MB per day means thousands of tiny shards within a year, which wastes heap and slows the cluster. Use rollover by size, fewer primaries, and ILM deletion.
Unbounded Indices
A single logs index that grows forever can't be tiered or deleted cheaply, and its shards become huge. Use data streams with rollover.
Treating Elasticsearch as the Source of Truth
Mapping changes require reindexing, and without an authoritative source you can't rebuild. Keep primary data in a database, or at least keep raw events in durable storage.
FAQ
How many shards should my index have?
Enough that each shard lands roughly in the 10–50 GB range for expected data, and no more. A few-GB index needs one primary shard. For time-series data, control shard size with rollover (max_primary_shard_size) rather than guessing counts up front.
What's the difference between primary shards and replicas?
Primary shards hold the original partitions of an index's data, and their number is fixed at creation. Replicas are copies of primaries on other nodes, providing failover and additional read capacity, and they can be added or removed anytime.
What does a yellow cluster status mean?
All primary shards are allocated, but at least one replica isn't, often because there aren't enough nodes to place replicas on different nodes than their primaries. Data is available, but you have less redundancy. On a single-node dev cluster with replicas configured, yellow is expected.
How do I reindex without downtime?
Create a new index with the updated mapping, reindex into it (or rebuild from the source database), keep it updated with ongoing changes, then atomically switch the alias your application uses from the old index to the new one.
Related Topics
- Elasticsearch — The search engine overview
- Elasticsearch Mappings — Templates and reindexing
- Database Sharding — Sharding concepts
- Log Aggregation — Time-series indices at scale
- Change Data Capture — Syncing from databases
- High Availability — Replicas and failover