NOTE
1.4 Elasticsearch Architecture
English translation of the original VNote ‘Elasticsearch Architecture’, preserving its structure and historical learning notes.
This is a historical learning note and may contain outdated or incomplete understanding.
1. Elasticsearch Cluster
1.1. A Cluster Is Composed of Nodes
- Nodes have two identities: master and data.
- A master node is responsible for managing nodes in the cluster, such as managing nodes and indices.
- A data node is responsible for storing document data and providing search.
- A node can have three status colors:
- red: a primary shard is not running normally.
- yellow: a replica shard is not running normally.
- green: primary shards and replica shards are running normally.
1.2. A Node Is Composed of Shards
There are two kinds: primary shards and replica shards. When writing a document, it is written to the primary shard and synchronized to replica shards. When reading a document, it can be read from either the primary shard or a replica shard.
The number of primary shards is fixed when the index is created, while the number of replica shards can be changed at any time. A primary shard cannot be placed on the same node as its own replica shard (otherwise, if the node goes down, both the primary shard and its replica are lost and fault tolerance is lost), but it can be placed on the same node as replica shards of other primary shards.
1.3. A Shard Is a Lucene Index
- In Elasticsearch, a shard is a Lucene Index, and a Lucene Index is composed of segments.
- A segment is essentially an inverted index and stores the actual documents.
1.3.1. Index, Type, Document
| Elasticsearch | Database | |
|---|---|---|
| Document | Row | |
| Type | Table | |
| Index | Database |
2. Elasticsearch Protocol
Similar to Raft, but not exactly the same.
2.1. Roles
- Leader: handles write requests.
- Follower: handles read requests.
2.2. Three Stages
2.2.1. Leader Election
- Three key points:
- Only candidate master nodes (nodes with
node.master: true) can become the master node. - The purpose of the minimum number of master nodes (
discovery.zen.minimum_master_nodes: (candidate master nodes / 2) + 1) is to prevent split brain. - Quorum mechanism: if more than half of the candidate master nodes vote for a candidate master node, that node is elected master.
- Only candidate master nodes (nodes with
- During each election, first filter the nodes that can become master (
node.master: true). - Then sort by nodeId lexicographically, and initially treat the first node as the master node.
- Aggregate the vote count. If the votes for a node reach a certain value (
n/2+1), that node is elected master. - Example:
- Suppose there are three nodes A, B, and C, all configured as:
node.master: true node.data: true discovery.zen.minimum_master_nodes: 2
- A starts. Through ping it gets the node list
(A)and elects the node with the smallest ID (only itself) as master, but the minimum of 2 nodes is not satisfied, so it waits in a loop. - B starts. Through ping it gets the node list
(A, B)and elects the node with the smallest ID, A, as master. The minimum of 2 nodes is satisfied, so they form a cluster. - C starts. Through ping it gets the node list
(A, B, C). Since there is already a master, it directly joins this cluster.
2.2.2. Request Processing
Elasticsearch CRUD Flow Elasticsearch Consistency
2.2.3. Crash Recovery
Same as Leader election.
3. Primary-Replica Replication
3.1. Leader Election
Select one replica from all replicas as Leader; the other replicas are Followers.
3.1.1. Election Method
- During each election, first filter the nodes that can become master (
node.master: true). - Then sort by nodeId lexicographically, and initially treat the first node as the master node.
- Aggregate the vote count. If the votes for a node reach a certain value (
n/2+1), that node is elected master. - Example:
- Suppose there are three nodes A, B, and C, all configured as:
node.master: true node.data: true discovery.zen.minimum_master_nodes: 2
- A starts. Through ping it gets the node list
(A)and elects the node with the smallest ID (only itself) as master, but the minimum of 2 nodes is not satisfied, so it waits in a loop. - B starts. Through ping it gets the node list
(A, B)and elects the node with the smallest ID, A, as master. The minimum of 2 nodes is satisfied, so they form a cluster. - C starts. Through ping it gets the node list
(A, B, C). Since there is already a master, it directly joins this cluster.
3.1.2. Split-Brain Problem
The purpose of the minimum number of master nodes (discovery.zen.minimum_master_nodes: (candidate master nodes / 2) + 1) is to prevent split brain.
3.2. Data Synchronization
3.2.1. Synchronization Process
- When a follower connects to the leader for the first time, it needs to synchronize all of the leader’s data. This process is called full synchronization. The process is:
- The leader takes a snapshot of the data at the current moment.
- The leader sends the snapshot to the new follower.
- The leader continues serving client writes.
- The follower replays the snapshot.
- The follower pulls all data changes that occurred after the leader snapshot.
3.2.2. Synchronization Method
Synchronous.
3.2.3. Synchronization Log
Elasticsearch does not use a log; instead, the primary shard sends synchronization requests in parallel to replica shards for synchronization.
3.3. Request Processing
3.3.1. Read Requests
Can be handled by the leader or a follower.
3.3.2. Write Requests
Must be handled by the leader. If the request is routed to the leader, the leader handles it and then synchronizes it to followers. If the request is routed to a follower, the follower must forward it to the leader; after the leader handles it, it synchronizes it to followers.
3.4. Failure Handling
3.4.1. Failure Detection
3.4.2. Failure Recovery
3.4.2.1. Follower Crash
- After a follower crashes and restarts, it can determine from its local log up to which position it has replicated. After reconnecting to the leader, it can continue replication from that position. This is called incremental synchronization. The process is:
- The follower restarts and connects to the leader.
- The follower reads the local log position and pulls data changes after that position from the leader.
- The follower replays those data changes.
3.4.2.2. Leader Crash
- After the leader crashes, the system needs to:
- Select a follower and promote it to leader.
- Notify clients and other followers that the leader has changed.
4. Partitioning
Distributed-System Partitioning
4.1. Split Data
4.1.1. Choosing the Split Key
Elasticsearch uses an internally generated ID as the key by default, and an ID can also be specified manually as the key.
4.1.2. Data-Splitting Strategy
hash.
4.2. Partition Allocation
Dynamic allocation. Automatic partition rebalancing.
4.3. Request Processing
4.3.1. Routing Component
server.
4.3.2. Routing Process
Calculate which shard the document should be on based on _id, i.e. hash(_id) % number_of_primary_shards, and then use the cluster state to determine which node hosts that shard.
Elasticsearch CRUD Flow
Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub