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 …
LeaderboardSalaryAccount