| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172 | 
[[docs-replication]]=== Reading and Writing documents[discrete]==== IntroductionEach index in Elasticsearch is <<scalability,divided into shards>>and each shard can have multiple copies. These copies are known as a _replication group_ and must be kept in sync when documentsare added or removed. If we fail to do so, reading from one copy will result in very different results than reading from another.The process of keeping the shard copies in sync and serving reads from them is what we call the _data replication model_.Elasticsearch’s data replication model is based on the _primary-backup model_ and is described very well in thehttps://www.microsoft.com/en-us/research/publication/pacifica-replication-in-log-based-distributed-storage-systems/[PacificA paper] ofMicrosoft Research. That model is based on having a single copy from the replication group that acts as the primary shard.The other copies are called _replica shards_. The primary serves as the main entry point for all indexing operations. It is in charge ofvalidating them and making sure they are correct. Once an index operation has been accepted by the primary, the primary is alsoresponsible for replicating the operation to the other copies.This purpose of this section is to give a high level overview of the Elasticsearch replication model and discuss the implicationsit has for various interactions between write and read operations.[discrete][[basic-write-model]]==== Basic write modelEvery indexing operation in Elasticsearch is first resolved to a replication group using <<index-routing,routing>>,typically based on the document ID. Once the replication group has been determined, the operation is forwardedinternally to the current _primary shard_ of the group. This stage of indexing is referred to as the _coordinating stage_.The next stage of indexing is the _primary stage_, performed on the primary shard. The primary shard is responsiblefor validating the operation and forwarding it to the other replicas. Since replicas can be offline, the primaryis not required to replicate to all replicas. Instead, Elasticsearch maintains a list of shard copies that shouldreceive the operation. This list is called the _in-sync copies_ and is maintained by the master node. As the name implies,these are the set of "good" shard copies that are guaranteed to have processed all of the index and delete operations thathave been acknowledged to the user. The primary is responsible for maintaining this invariant and thus has to replicate alloperations to each copy in this set.The primary shard follows this basic flow:. Validate incoming operation and reject it if structurally invalid (Example: have an object field where a number is expected). Execute the operation locally i.e. indexing or deleting the relevant document. This will also validate the content of fields   and reject if needed (Example: a keyword value is too long for indexing in Lucene).. Forward the operation to each replica in the current in-sync copies set. If there are multiple replicas, this is done in parallel.. Once all replicas have successfully performed the operation and responded to the primary, the primary acknowledges the successful   completion of the request to the client.Each in-sync replica copy performs the indexing operation locally so that it has a copy. This stage of indexing is the_replica stage_.These indexing stages (coordinating, primary, and replica) are sequential. To enable internal retries, the lifetime of each stageencompasses the lifetime of each subsequent stage. For example, the coordinating stage is not complete until each primarystage, which may be spread out across different primary shards, has completed. Each primary stage will not complete until thein-sync replicas have finished indexing the docs locally and responded to the replica requests.[discrete]===== Failure handlingMany things can go wrong during indexing -- disks can get corrupted, nodes can be disconnected from each other, or someconfiguration mistake could cause an operation to fail on a replica despite it being successful on the primary. Theseare infrequent but the primary has to respond to them.In the case that the primary itself fails, the node hosting the primary will send a message to the master about it. The indexingoperation will wait (up to 1 minute, by <<dynamic-index-settings,default>>) for the master to promote one of the replicas to be anew primary. The operation will then be forwarded to the new primary for processing. Note that the master also monitors thehealth of the nodes and may decide to proactively demote a primary. This typically happens when the node holding the primaryis isolated from the cluster by a networking issue. See <<demoted-primary,here>> for more details.Once the operation has been successfully performed on the primary, the primary has to deal with potential failureswhen executing it on the replica shards. This may be caused by an actual failure on the replica or due to a networkissue preventing the operation from reaching the replica (or preventing the replica from responding). All of theseshare the same end result: a replica which is part of the in-sync replica set misses an operation that is about tobe acknowledged. In order to avoid violating the invariant, the primary sends a message to the master requestingthat the problematic shard be removed from the in-sync replica set. Only once removal of the shard has been acknowledgedby the master does the primary acknowledge the operation. Note that the master will also instruct another node to startbuilding a new shard copy in order to restore the system to a healthy state.[[demoted-primary]]While forwarding an operation to the replicas, the primary will use the replicas to validate that it is still theactive primary. If the primary has been isolated due to a network partition (or a long GC) it may continue to processincoming indexing operations before realising that it has been demoted. Operations that come from a stale primarywill be rejected by the replicas. When the primary receives a response from the replica rejecting its request becauseit is no longer the primary then it will reach out to the master and will learn that it has been replaced. Theoperation is then routed to the new primary..What happens if there are no replicas?************This is a valid scenario that can happen due to index configuration or simplybecause all the replicas have failed. In that case the primary is processing operations without any external validation,which may seem problematic. On the other hand, the primary cannot fail other shards on its own but request the master to doso on its behalf. This means that the master knows that the primary is the only single good copy. We are therefore guaranteedthat the master will not promote any other (out-of-date) shard copy to be a  new primary and that any operation indexedinto the primary will not be lost. Of course, since at that point we are running with only single copy of the data, physical hardwareissues can cause data loss. See <<index-wait-for-active-shards>> for some mitigation options.************[discrete]==== Basic read modelReads in Elasticsearch can be very lightweight lookups by ID or a heavy search request with complex aggregations thattake non-trivial CPU power. One of the beauties of the primary-backup model is that it keeps all shard copies identical(with the exception of in-flight operations). As such, a single in-sync copy is sufficient to serve read requests.When a read request is received by a node, that node is responsible for forwarding it to the nodes that hold the relevant shards,collating the responses, and responding to the client. We call that node the _coordinating node_ for that request. The basic flowis as follows:. Resolve the read requests to the relevant shards. Note that since most searches will be sent to one or more indices,   they typically need to read from multiple shards, each representing a different subset of the data.. Select an active copy of each relevant shard, from the shard replication group. This can be either the primary or   a replica. By default, Elasticsearch will simply round robin between the shard copies.. Send shard level read requests to the selected copies.. Combine the results and respond. Note that in the case of get by ID look up, only one shard is relevant and this step can be skipped.[discrete][[shard-failures]]===== Shard failuresWhen a shard fails to respond to a read request, the coordinating node sends therequest to another shard copy in the same replication group. Repeated failurescan result in no available shard copies.To ensure fast responses, the following APIs willrespond with partial results if one or more shards fail:* <<search-search, Search>>* <<search-multi-search, Multi Search>>* <<docs-bulk, Bulk>>* <<docs-multi-get, Multi Get>>Responses containing partial results still provide a `200 OK` HTTP status code.Shard failures are indicated by the `timed_out` and `_shards` fields ofthe response header.[discrete]==== A few simple implicationsEach of these basic flows determines how Elasticsearch behaves as a system for both reads and writes. Furthermore, since readand write requests can be executed concurrently, these two basic flows interact with each other. This has a few inherent implications:Efficient reads:: Under normal operation each read operation is performed once for each relevant replication group.   Only under failure conditions do multiple copies of the same shard execute the same search.Read unacknowledged:: Since the primary first indexes locally and then replicates the request, it is possible for a  concurrent read to already see the change before it has been acknowledged.Two copies by default:: This model can be fault tolerant while maintaining only two copies of the data. This is in contrast to  quorum-based system where the minimum number of copies for fault tolerance is 3.[discrete]==== FailuresUnder failures, the following is possible:A single shard can slow down indexing:: Because the primary waits for all replicas in the in-sync copies set during each operation,  a single slow shard can slow down the entire replication group. This is the price we pay for the read efficiency mentioned above.  Of course a single slow shard will also slow down unlucky searches that have been routed to it.Dirty reads:: An isolated primary can expose writes that will not be acknowledged. This is caused by the fact that an isolated  primary will only realize that it is isolated once it sends requests to its replicas or when reaching out to the master.  At that point the operation is already indexed into the primary and can be read by a concurrent read. Elasticsearch mitigates  this risk by pinging the master every second (by default) and rejecting indexing operations if no master is known.[discrete]==== The Tip of the IcebergThis document provides a high level overview of how Elasticsearch deals with data. Of course, there is much much moregoing on under the hood. Things like primary terms, cluster state publishing, and master election all play a role inkeeping this system behaving correctly. This document also doesn't cover known and importantbugs (both closed and open). We recognize that https://github.com/elastic/elasticsearch/issues?q=label%3Aresiliency[GitHub is hard to keep up with].To help people stay on top of those, we maintain a dedicated https://www.elastic.co/guide/en/elasticsearch/resiliency/current/index.html[resiliency page]on our website. We strongly advise reading it.
 |