Design an Elasticsearch Cluster
Problem Given a target number of active shards, a number of nodes, and a required replica count, design an Elasticsearch cluster topology and explain how the data is stored and served.
Functional requirements
- Lay out primary and replica shards across the available nodes.
- Survive the loss of a node without data loss or write unavailability.
- Serve both indexing and search traffic from the same cluster.
- Allow the index to grow without a full reindex where possible.
- Support rolling restarts/upgrades with no downtime.
Non-functional requirements
- ~2 TB primary data, targeting 30-50 GB per shard → ~50 primary shards; with replicas=1 that is ~100 shard copies.
- Sized across ~10 data nodes → ~10 shard copies/node, well under the ~20 shards per GB of heap rule of thumb (31 GB heap cap).
- ~5k index ops/sec sustained, ~2k search QPS at peak.
- Search p99 < 200 ms; near-real-time visibility within the 1 s default refresh interval.
- Tolerate 1 node loss with zero data loss; 3 dedicated master-eligible nodes for quorum.
Key components
- Primary shards: fixed at index creation, cap write parallelism and determine data distribution. Changing the count requires reindex or the split/shrink APIs.
- Replica shards: copies for fault tolerance and read scaling; a replica is never allocated to the same node as its primary, so replicas > nodes-1 leaves shards permanently unassigned (yellow health).
- Node roles: dedicated master-eligible nodes (3, for quorum and split-brain avoidance), data nodes, and coordinating nodes fronting client traffic.
- Storage internals: each shard is a self-contained Lucene index of immutable segments built over an inverted index; documents are analyzed/tokenized at index time; refresh makes them searchable, flush/translog makes them durable, and background merges compact segments.
- Allocation awareness across racks/AZs so a rack loss cannot take out both a primary and its replica.
Deep dives / trade-offs
- Shard sizing: over-sharding is the classic failure — each shard carries file handles, heap for segment metadata, and cluster-state bloat, so thousands of tiny shards will collapse a cluster. Under-sharding caps write parallelism and produces huge, slow-merging shards. Work the arithmetic from data volume and the 30-50 GB target.
- Replicas trade write cost for read throughput and safety: each replica multiplies indexing work and disk, but adds a search-serving copy. Discuss setting replicas=0 during a bulk backfill and raising it after.
- Failure handling: losing a node with a primary promotes an in-sync replica automatically; walk through what happens with replicas=0 (data loss), and what the translog guarantees for in-flight writes.
- Refresh vs flush vs merge: raising refresh_interval during bulk loads dramatically improves throughput at the cost of search visibility latency.
- Hot/warm architecture and time-based indices with ILM when the data is append-mostly.
asked …