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=geektimeat startup.
- Clusters are distinguished by unique names, defaulting to
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_nodesto 2.
Starting from version 7.0, this configuration is no longer needed:
- The
minimum_master_nodesparameter 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
_routingvalue 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
- The coordinating node sends the query to all shards involved.
- Each shard runs the query locally and returns a sorted list of document IDs and scores to the coordinating node.
- The coordinating node merges the results and sorts them globally.
Fetch Phase
- The coordinating node requests the actual document contents from the relevant shards.
- The shards return the documents, and the coordinating node assembles the final result.
Issues with Query-Then-Fetch
Performance:
- Each shard retrieves
from + sizedocuments. - 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
1primary shard. - For large datasets, ensure documents are evenly distributed across shards.
- Use
DFS Query Then Fetchby adding?search_type=dfs_query_then_fetchto 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
textfields). - Fielddata (for
textfields, 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
textfields. - Sorting on
textfields 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/sizefor straightforward pagination,search_afterfor 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
}