Elasticsearch Distributed Architecture and Search Mechanics

Cluster State Maintenance, Leader Election, and Split-Brain

Distributed Characteristics

  • Benefits of distributed architecture in Elasticsearch:
    • Horizontal scaling of storage, supporting petabyte-scale data.
    • Increased system availability; failure of some nodes does not affect the entire cluster.
  • Distributed architecture details:
    • Clusters are distinguished by unique names, defaulting to elasticsearch.
    • The cluster name can be set in the configuration file or via -E cluster.name=geektime at startup.

Nodes

A node is an instance of Elasticsearch, essentially a Java process. Although multiple Elasticsearch processes can run on a single machine, production environments typically recommend running one instance per machine.

Each node has a name, set via configuration file or -E node.name=geektime at startup. After starting, each node is assigned a UID stored in the data directory.

Coordinating Node

A node that handles client requests is called a coordinating node. It routes requests to the appropriate nodes, e.g., an index creation request is routed to the Master node.

All nodes are coordinating nodes by default. To create a dedicated coordinating node, set other node roles to false.

Demo – Starting Nodes and Cerebro Overview

Start a cluster from the command line:

bin/elasticsearch -E node.name=node1 -E cluster.name=geektime -E path.data=node1_data -E http.port=9200 -E network.host=0.0.0.0 -E node.master=true
bin/elasticsearch -E node.name=node2 -E cluster.name=geektime -E path.data=node2_data -E http.port=9201 -E network.host=0.0.0.0 -E node.master=true
bin/elasticsearch -E node.name=node3 -E cluster.name=geektime -E path.data=node3_data -E http.port=9202 -E network.host=0.0.0.0 -E node.master=true

Download Cerebro from https://github.com/lmenezes/cerebro/releases.

Access Cerebro at http://192.168.163.131:9000/.

Create an index named test with 3 primary shards and 1 replica.

Data Node

Nodes that can store data are called data nodes. By default, all nodes are data nodes; set node.data: false to disable this role.

Data node responsibilities:

  • Store shard data.
  • Enable data scaling (the Master node decides how shards are distributed among data nodes).

Adding data nodes solves horizontal scaling and single points of failure for data.

Master Node

Master node responsibilities:

  • Handle index creation, deletion, and other administrative requests.
  • Decide which nodes host shards.
  • Maintain and update the cluster state.

Best practices for Master nodes:

  • Master nodes are critical; avoid single points of failure.
  • Configure multiple Master-eligible nodes, each dedicated to the Master role.

Master-Eligible Nodes and Leader Election

A cluster can have multiple Master-eligible nodes. These nodes can participate in leader election, becoming the Master node when necessary (e.g., if the current Master fails or a network partition occurs).

Each node is Master-eligible by default; set node.master: false to disable this role.

When the first Master-eligible node starts, it elects itself as the Master node.

Cluster State

The cluster state contains essential information:

  • All node information.
  • All indices, their mappings, and settings.
  • Shard routing information.

Every node maintains a copy of the cluster state, but only the Master node can modify it and synchronize changes to other nodes. Allowing any node to modify the state would cause inconsistency.

Leader Election Process

  • Master-eligible nodes ping eachother. The node with the lowest node ID is elected as the Master.
  • Other nodes join the cluster but do not take on the Master role. If the elected Master is lost, a new election occurs.

Split-Brain Problem

Split-brain is a classic network issue in distributed systems where a node becomes disconnected from the rest of the cluster.

Example:

  • Node 2 and Node 3 lose connection with Node 1.
  • Node 2 and Node 3 elect a new Master.
  • Node 1 remains Master, forming a separate cluster and updating its cluster state.
  • This results in two Masters maintaining different cluster states. When the network recovers, it is difficult to determine the correct state.

Avoiding Split-Brain

To prevent split-brain, set a quorum for elections. An election can proceed only if the number of Master-eligible nodes exceeds the quorum.

  • Quorum = (total Master-eligible nodes / 2) + 1
  • With 3 Master-eligible nodes, set discovery.zen.minimum_master_nodes to 2.

Starting from version 7.0, this configuration is no longer needed:

  • The minimum_master_nodes parameter is removed; Elasticsearch selects nodes that can form a quorum.
  • Leader elections complete quickly.
  • Cluster scaling is safer and easier, with fewer system configuration options that could cause data loss.
  • Nodes log their state more clearly, aiding in diagnosing why they cannot join the cluster or elect a Master.

Configuring Node Types

By default, a node is Master-eligible, a data node, and an ingest node.

Shards and Cluster Failover

Primary Shards – Increasing Storage Capacity

Shards are the foundation of Elasticsearch distributed storage. They are divided into primary shards and replica shards.

Primary shards distribute index data across multiple data nodes, enabling horizontal scaling. The number of primary shards is set when creating an index and cannot be changed later without reindexing.

Replica Shards – Enhancing Data Availability

Replica shards improve data availability. If a primary shard is lost, a replica can be promoted to primary. The number of replicas can be adjusted dynamically. Without replicas, hardware failure can lead to data loss.

Replica shards also improve read performance. They are synchronized from the primary shard. Increasing the number of replicas can increase read throughput.

Shard Count Planning

When planning the number of primary and replica shards:

  • Too few primary shards: A rapidly growing index cannot be scaled by adding nodes.
  • Too many primary shards: Creates many small shards, potentially harming performance.
  • Too many replica shards: Decreases overall write performance.

Cluster Health Status

The cluster health status is indicated by green (all shards active), yellow (only primary shards active), or red (some primary shards missing).

Document Distributed Storage

Document-to-Shard Mapping

Each document is stored on a specific primary shard and its replicas. For example, document 1 is stored on shards P0 and R0.

The algorithm maps documents to shards to ensure even distribution across all shards, optimizing hardware usage.

Routing Algorithm

shard = hash(_routing) % number_of_primary_shards
  • The hash function ensures documents are evenly distributed.
  • The default _routing value is the document ID.
  • Custom routing values can group related documents on the same shard (e.g., products from the same country).
  • This explains why the number of primary shards cannot be easily changed after index creation.

Updating and Deleting Documents

When updating a document, Elasticsearch marks the old document as deleted and indexes a new one. The document version is incremented. Deletion marks the document as deleted but does not immediately free disk space.

Shard Lifecycle

Internal Shard Mechanics

A shard is a Lucene index, which is the smallest unit of work in Elasticsearch.

Key concepts:

  • Why Elasticsearch is near real-time (documents become searchable after about 1 second).
  • How Elasticsearch ensures no data loss during power failures.
  • Why deleting documents does not immediately reclaim space.

Inverted Index Immutability

The inverted index uses an immutable design—once created, it cannot be changed.

Benefits:

  • No need for concurrency control on file writes, avoiding locks.
  • Once loaded into the file system cache, it stays there. If enough cache is available, most requests are served from memory, improving performance.
  • Caching is easy to generate and maintain; data can be compressed.

Challenges:

  • To make a new document searchable, the entire index must be rebuilt.

Lucene Index Operations

  • Refresh: Makes newly indexed documents visible to search. This is why search is near real-time.
  • Transaction Log (translog): Records all operations to ensure durability. Used to recover data after a crash.
  • Flush: Forces a commit, persisting segments to disk and clearing the translog.

Segment Merging

Over time, many small segments accumulate. Elasticsearch and Lucene automatically merge them to:

  • Reduce the number of segments.
  • Permanently remove deleted documents.

Manual merge can be triggered via POST my_index/_forcemerge.

Distributed Search and Relevance Scoring

Query-Then-Fetch Mechanism

Elasticsearch performs search in two phases: Query and Fetch.

Query Phase

  1. The coordinating node sends the query to all shards involved.
  2. Each shard runs the query locally and returns a sorted list of document IDs and scores to the coordinating node.
  3. The coordinating node merges the results and sorts them globally.

Fetch Phase

  1. The coordinating node requests the actual document contents from the relevant shards.
  2. The shards return the documents, and the coordinating node assembles the final result.

Issues with Query-Then-Fetch

Performance:

  • Each shard retrieves from + size documents.
  • The coordinating node processes number_of_shard * (from + size) documents.
  • Deep pagination is problematic.

Relevance Scoring:

  • Each shard computes scores based on its own data, leading to score deviations, especially with small datasets and multiple primary shards.

Solutions for Scoring Inaccuracies

  • For small datasets, use 1 primary shard.
  • For large datasets, ensure documents are evenly distributed across shards.
  • Use DFS Query Then Fetch by adding ?search_type=dfs_query_then_fetch to the search URL. This collects term and document frequencies from all shards before scoring, but it consumes more CPU and memory.
POST message/_search?search_type=dfs_query_then_fetch
{
  "query": {
    "term": {
      "content": {
        "value": "good"
      }
    }
  }
}

Sorting, Doc Values, and Fielddata

Sorting

By default, results are sorted by relevance score in descending order. Custom sorting can be specified using the sort parameter. If _score is not specified, it is not computed.

POST /kibana_sample_data_ecommerce/_search
{
  "size": 5,
  "query": {
    "match_all": {}
  },
  "sort": [
    {
      "order_date": {
        "order": "desc"
      }
    }
  ]
}

Multiple sort criteria can be combined. The first sort field takes priority.

POST /kibana_sample_data_ecommerce/_search
{
  "size": 5,
  "query": {
    "match_all": {}
  },
  "sort": [
    {
      "order_date": {
        "order": "desc"
      }
    },
    {
      "_doc": {
        "order": "asc"
      }
    },
    {
      "_score": {
        "order": "desc"
      }
    }
  ]
}

Sorting Process

Sorting operates on the original field values, not the inverted index. It uses a forward index, which can be implemented via:

  • Doc Values (default, columnar storage, not available for text fields).
  • Fielddata (for text fields, off by default).

Doc Values vs. Fielddata

Fielddata:

  • Off by default; can be enabled via mapping. Changes take effect immediately without reindexing.
  • Only available for text fields.
  • Sorting on text fields is possible but results may not be useful; not recommended.
  • Can be used for certain aggregation needs.

Enabling Fielddata:

PUT my_index/_mapping
{
  "properties": {
    "my_field": {
      "type": "text",
      "fielddata": true
    }
  }
}

Doc Values:

  • Enabled by default; can be disabled via mapping to improve indexing speed and reduce disk space.
  • If disabled, re-enabling requires reindexing.
  • Disable when sorting and aggregation are not needed.

Disabling Doc Values:

PUT my_index/_mapping
{
  "properties": {
    "my_field": {
      "type": "keyword",
      "doc_values": false
    }
  }
}

Viewing Doc Values and Fielddata Content

text fields do not support Doc Values. With Fielddata enabled for text fields, stored tokenized data can be viewed.

DELETE temp_users
PUT temp_users
PUT temp_users/_mapping
{
  "properties": {
    "name": {"type": "text","fielddata": true},
    "desc": {"type": "text","fielddata": true}
  }
}

POST temp_users/_doc
{"name":"Jack","desc":"Jack is a good boy!","age":10}

POST temp_users/_search
{
  "docvalue_fields": [
    "name",
    "desc"
  ]
}

For numeric fields, Doc Values are used:

POST temp_users/_search
{
  "docvalue_fields": [
    "age"
  ]
}

Pagination and Iteration: From/Size, Search After, and Scroll API

From/Size

By default, queries return the first 10 results sorted by score. Pagination uses:

  • from: Start position.
  • size: Number of documents to return.

Deep Pagination in Distributed Systems

from + size must be less than 10,000 by default. Deep pagination is inefficient.

Search After

search_after avoids the performance cost of deep pagination and provides a live cursor for subsequent pages.

  • Does not support specifying a page number (from).
  • Only moves forward.
  • Requires a sort with unique values (e.g., include _id).
  • Uses the sort values of the last document from the previous page.

Example data ingestion:

DELETE users
POST users/_doc
{"name":"user1","age":10}
POST users/_doc
{"name":"user2","age":11}
POST users/_doc
{"name":"user2","age":12}
POST users/_doc
{"name":"user2","age":13}
POST users/_count

Initial search with sort:

POST users/_search
{
    "size": 1,
    "query": {
        "match_all": {}
    },
    "sort": [
        {"age": "desc"} ,
        {"_id": "asc"}
    ]
}

Subsequent page using search_after:

POST users/_search
{
    "size": 1,
    "query": {
        "match_all": {}
    },
    "search_after": [13, "cuDPj3gBXBR-8Kp6kAmc"],
    "sort": [
        {"age": "desc"} ,
        {"_id": "asc"}
    ]
}

search_after solves deep pagination by maintaining a cursor instead of skipping results.

Scroll API

Creates a snapshot of the data at a specific point in time. New writes after the snapshot are not visible.

Each query returns a scroll_id to use in subsequent requests.

Creating a snapshot with a 5-minute lifetime:

DELETE users
POST users/_doc
{"name":"user1","age":10}
POST users/_doc
{"name":"user2","age":20}
POST users/_doc
{"name":"user3","age":30}
POST users/_doc
{"name":"user4","age":40}

POST /users/_search?scroll=5m
{
  "size": 1,
  "query": {
    "match_all": {}
  }
}

Using the scroll to get the next batch:

POST /_search/scroll
{
  "scroll": "1m",
  "scroll_id": "DXF1ZXJ5QW5kRmV0Y2gBAAAAAAAAA5EWT0VUdC1XNThUUHVWM1pzS0NSck9Zdw=="
}

Use Cases for Different Search Types

  • Regular search (default): Top results, e.g., latest orders.
  • Scroll: All documents, e.g., exporting data.
  • Pagination: from/size for straightforward pagination, search_after for deep or real-time pagination.

Handling Concurrent Read/Write Operations

Optimistic Concurrency Control

Documents in Elasticsearch are immutable. Updates mark the old document as deleted and index a new one, incrementing the document version.

Internal version control uses:

  • if_seq_no + if_primary_term

External version control (e.g., using another database):

  • version + version_type=external

Example:

DELETE products
PUT products

PUT products/_doc/1
{
  "title": "iphone",
  "count": 100
}

GET products/_doc/1

Update with optimistic locking by specifying seq_no and primary_term:

PUT products/_doc/1?if_seq_no=0&if_primary_term=1
{
  "title": "iphone",
  "count": 100
}

Using a custom version:

PUT products/_doc/1?if_seq_no=100&if_primary_term=1
{
  "title": "iphone",
  "count": 100
}

PUT products/_doc/1?version=100&version_type=external
{
  "title": "apple",
  "count": 100
}

PUT products/_doc/1?if_seq_no=101&if_primary_term=1
{
  "title": "app",
  "count": 10
}

Tags: elasticsearch Distributed Systems Cluster Architecture Sharding Replication

Posted on Fri, 02 Oct 2026 16:28:16 +0000 by Eclectic