Lesson 3 / 25

Cluster, Node, Index, Shard and Replica

The building blocks.

How data is spread out

A cluster is a group of nodes (Elasticsearch processes) sharing a cluster name; one elected master node manages cluster state such as mappings and shard locations. Data is stored in indices, each a logical collection of documents (JSON objects with an _id). An index is split into one or more primary shards, each a self-contained Lucene index, and every primary can have replica shards: copies on other nodes that provide failover and extra read capacity. Writes go to the primary and are then replicated; searches can be served by any copy. Since 7.0 a new index defaults to one primary and one replica. The number of primaries is fixed at creation (changing it needs split, shrink or reindex); the number of replicas can be changed at any time.

Creating an index with explicit shard settings

Kibana Dev Tools console syntax; send the same requests with curl or a client library.

PUT /orders
{
  "settings": {
    "number_of_shards": 3,
    "number_of_replicas": 1
  }
}

# replicas can be changed later; primaries cannot
PUT /orders/_settings
{
  "index": { "number_of_replicas": 2 }
}

# see where each shard copy lives
GET /_cat/shards/orders?v

Chapters and photocopies

An index is a book split into chapters (primary shards) handed to different people (nodes). Each chapter also has a photocopy held by someone else (a replica), so losing one person loses no chapter.

Quick check: Which index setting can be changed on a live index without reindexing?

  • The index name
  • number_of_shards
  • The type of an existing field
  • number_of_replicas
Answer

number_of_replicas — Replica count is dynamic; primary count is fixed at creation.