NOTE

1.3 Elasticsearch CRUD Flow

English translation of the original VNote ‘Elasticsearch CRUD Flow’, preserving its steps, examples, code, and historical questions.

Elasticsearch / SearchCreated Updated 4 min readhistorical

This is a historical learning note and may contain outdated or incomplete understanding.

1. Writing Data

  1. The client randomly selects a node and sends a write request.
  2. This node acts as the coordinating node. It calculates which shard the document should be on based on _id, i.e. hash(_id) % number_of_primary_shards, and then obtains which node hosts that shard from the cluster state.
  3. Route the request to the primary shard on the corresponding node.
  4. The primary shard performs the write.
    1. The primary shard writes the data into the memory buffer.
    2. The primary shard writes the data into the transaction log.
  5. The primary shard sends the data to replica shards in parallel.
  6. After synchronization completes, the replica shards respond success to the primary shard.
  7. The primary shard returns success to the coordinating node.
  8. The coordinating node returns success to the client.
  • ES write data - ES write flow

2. Deleting Data

2.1. Delete by ID

  1. The client randomly selects a node and sends an ID-delete request.
  2. This node acts as the coordinating node. It calculates which shard the document should be on based on _id, i.e. hash(_id) % number_of_primary_shards, and then obtains which node hosts that shard from the cluster state.
  3. The coordinating node routes the request to the primary shard.
  4. The primary shard performs the delete operation.
    1. The primary shard writes the change into the memory buffer.
      • At commit time, this record is written into a .del file, indicating that the record has been deleted from the segment. At this point the document can still match queries, but it will be filtered from results.
      • Elasticsearch refresh
    2. The primary shard writes the operation into the transaction log.
      • During flush, segment merging is performed, and data listed in the .del file will not be written into the new segment.
      • Elasticsearch translog
  5. The primary shard sends the delete request to replica shards in parallel.
  6. After synchronization completes, the replica shards respond success to the primary shard.
  7. The primary shard returns success to the coordinating node.
  8. The coordinating node returns success to the client.

2.2. delete_by_query

  1. The client randomly selects a node and sends a delete_by_query request.
  2. This node acts as the coordinating node and routes the delete_by_query request to all primary shards.
  3. Each primary shard performs the delete operation.
    1. First, find the relevant documents (mainly including version + id).
    2. Then perform delete-by-ID while comparing the version in the translog (question here: Elasticsearch keeps tracks of the sequence number and primary term of the last operation to have changed each of the documents it stores.). If they do not match, report a version conflict; otherwise deletion succeeds.
  4. Each primary shard sends the delete operation to replica shards in parallel.
  5. Replica shards return delete success to the primary shard.
  6. Primary shards return delete success to the coordinating node.
  7. The coordinating node returns delete success to the client.
  • Example
    // Disable automatic refresh
    PUT /tb_item/_settings
    {
      "index" : {
        "refresh_interval" : -1
      }
    }
    // Delete
    POST /tb_item/_delete_by_query
    {
      "query": {
        "match": {
          "id": "936920"
        }
      }
    }
    // It can still be found by search
    GET /tb_item/_search
    {
      "query": {
        "ids" : {
          "values" : ["936920"]
        }
      }
    }
    // Updating by ID reports that the document does not exist
    POST /tb_item/_update/936920
    {
      "doc": {
        "sellPoint": "test1"
      }
    }

3. Querying Data

3.1. Query by ID

  1. The client randomly selects a node and sends a read request.
  2. This node acts as the coordinating node. It calculates which shard the document should be on based on _id, i.e. hash(_id) % number_of_primary_shards, and then obtains which node hosts that shard from the cluster state.
  3. The coordinating node routes the request to either the primary shard or a replica shard.
  4. The shard obtains the document and returns it to the coordinating node.
  5. The coordinating node returns the data to the client.

3.2. Keyword Query

  1. The client randomly selects a node and sends a read request.
  2. This node acts as the coordinating node and routes the read request to all shards (either primary or replica shards can be used).
  3. Each shard returns key information about the documents it found (including _id) to the coordinating node. The coordinating node merges, sorts, paginates, and otherwise processes the data to produce the final result. This operation is called the query phase.
  4. The coordinating node takes the final _id values and fetches the actual document data from the corresponding nodes, then returns it to the client. This operation is called the fetch phase.

Why split this into two phases instead of one? For example, why not have every shard return all document information directly to the coordinating node? Because returning all of that data would be too large.

3.2.1. Paginated Query

  1. The client randomly selects a node and sends a read request to get 10 records.
  2. This node acts as the coordinating node and routes the read request to all shards (either primary or replica shards can be used).
  3. Each shard returns key information about its own top 10 documents (including _id) to the coordinating node. The coordinating node merges, sorts, paginates, and otherwise processes the data to produce the final 10 results. This is the query phase.
  4. The coordinating node takes the final _id values and fetches the actual document data from the corresponding nodes, then returns it to the client. This is the fetch phase.

4. Updating Data

4.1. Update by ID

  1. Delete.
    1. Refer to Delete by ID.
  2. Write.
    1. Refer to Writing Data, version + 1.

4.2. update_by_query

  1. The client randomly selects a node and sends an update_by_query request.
  2. This node acts as the coordinating node and routes the read request to all shards (either primary or replica shards can be used).
  3. Each primary shard performs the update operation.
    1. First, find the relevant documents (mainly including version + id).
    2. Then perform update-by-ID while comparing the version in the translog (question here: Elasticsearch keeps tracks of the sequence number and primary term of the last operation to have changed each of the documents it stores.). If they do not match, report a version conflict; otherwise the update succeeds.
  4. Each primary shard sends the update operation to replica shards in parallel.
  5. Replica shards return update success to the primary shard.
  6. Primary shards return update success to the coordinating node.
  7. The coordinating node returns update success to the client.

5. References

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub