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

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

Shard Sizing

Each shard has overhead (heap, file handles, cluster state), so too many small shards (oversharding) is the most common scaling mistake. Guidelines:

See database sharding.

Cluster Health

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):

This keeps shards well-sized and costs proportional to data value. See log aggregation.

Indexing Throughput

Keeping Elasticsearch in Sync

Elasticsearch is usually a secondary index, not the system of record. Sync patterns:

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

References