| = Distributed Update Internals (SolrCloud) |
| :toc: macro |
| :toclevels: 3 |
| |
| toc::[] |
| |
| == Why this doc |
| |
| The Solr Reference Guide states the user-facing consistency model of SolrCloud updates in |
| https://github.com/apache/solr/blob/main/solr/solr-ref-guide/modules/deployment-guide/pages/solrcloud-update-consistency.adoc[SolrCloud Update Consistency Model]. |
| This document explains *how* those guarantees are implemented: the distributed update path from the receiving node through the shard leader to the replicas, and the versioning scheme that everything else leans on. |
| |
| Class names below are under `solr/core/src/java/org/apache/solr/` unless otherwise noted. |
| |
| == Request lifecycle |
| |
| An `/update` request is parsed by a loader into a stream of `AddUpdateCommand` / `DeleteUpdateCommand` / `CommitUpdateCommand` objects, each fed through the update request processor (URP) chain. |
| Everything discussed here happens inside two processors near the end of that chain: `update.processor.DistributedZkUpdateProcessor` (subclass of `DistributedUpdateProcessor`; "DUP" below) and `update.processor.RunUpdateProcessor`. |
| URPs configured *before* the DUP run only on the node that received the client request; URPs *after* it (and `RunUpdateProcessor`) run on the leader and again on every replica. |
| |
| [source,mermaid] |
| ---- |
| sequenceDiagram |
| participant C as Client |
| participant N as Receiving node |
| participant L as Shard leader |
| participant R as NRT/TLOG replicas |
| |
| C->>N: /update (batch of docs) |
| N->>L: forward each doc (DistribPhase=TOLEADER) |
| L->>L: per-doc lock, OCC check, assign _version_ |
| L->>L: Lucene write + tlog append |
| L--)R: stream doc (DistribPhase=FROMLEADER, async) |
| R->>R: drop if stale version, else apply + tlog |
| L->>L: finish(): await replica responses,<br/>demote failed replicas via shard terms,<br/>flush tlog |
| L-->>N: per-request response (+ achieved rf) |
| N-->>C: HTTP 200 (or error) |
| ---- |
| |
| === Routing: setupRequest |
| |
| `DistributedZkUpdateProcessor.setupRequest()` decides, per document, what role this node plays. |
| It computes the target slice via the collection's `DocRouter`, looks up the shard leader with `ZkStateReader.getLeaderRetry(...)`, and compares it to the local core. |
| The `DISTRIB_UPDATE_PARAM` (`update.distrib`) carries a `DistribPhase` marking where the request came from: |
| |
| * `NONE` — an external client request; if this node is not the leader, set `forwardToLeader` and target a single `SolrCmdDistributor.ForwardNode` (the leader). |
| * `TOLEADER` — forwarded from a peer; this node is (should be) the leader; compute the replica fan-out list. |
| * `FROMLEADER` — forwarded from the leader; apply locally only, no further distribution. |
| |
| The fan-out list comes from `getReplicaNodesForLeader(...)`: only `NRT` and `TLOG` replicas (PULL replicas never receive updates), excluding down/non-live replicas and any replica whose shard term is already behind the leader's (`ZkShardTerms.skipSendingUpdatesTo`) — those are collected in `skippedCoreNodeNames` and will have their terms pushed further down at request end. |
| |
| Retry limits for internal hops are asymmetric: forwarding to a leader retries up to `solr.retries.on.forward` (default 25) times, while a leader sending to followers retries only `solr.retries.to.followers` (default 3) times — a follower that can't take the update is demoted instead of retried hard. |
| |
| === Leader-side processing |
| |
| For an add, `DistributedUpdateProcessor.versionAdd` runs the whole per-document decision inside `UpdateLocks.runWithLock(id, ...)` — see <<UpdateLocks>>. |
| `leaderLogic` is true when this core is the leader and the command is not a replay/peer-sync (`leaderLogicWithVersionIntegrityCheck`); a non-leader receiving an update *without* a version is rejected as an invalid state. |
| |
| On the leader: |
| |
| 1. If the command carries a client-supplied version constraint, perform the optimistic-concurrency check (see <<Leader-side optimistic concurrency>>). |
| 2. If the command is an atomic (partial) update, resolve it into a full replacement document (`getUpdatedDocument` → `AtomicUpdateDocumentMerger`), reading the current document through `RealTimeGetComponent.getInputDocument` — which sees uncommitted state via the tlog. |
| 3. Assign a new version: `cmd.setVersion(vinfo.getNewClock())` (see <<The _version_ clock>>). |
| 4. Apply locally: `doLocalAdd` → `RunUpdateProcessor` → `DirectUpdateHandler2.addDoc` (Lucene `updateDocument`) and `UpdateLog.add` (tlog append). |
| 5. Hand the (now fully-resolved, versioned) document to `doDistribAdd` for fan-out. |
| |
| Because steps 1–4 happen under the per-document lock, concurrent updates to the same id are serialized at the leader, read-modify-write atomic updates are linearizable per document, and version assignment order matches apply order per document. |
| |
| === Fan-out: SolrCmdDistributor and StreamingSolrClients |
| |
| `update.SolrCmdDistributor.distribAdd/distribDelete` submits each command to the target nodes. |
| Ordinary adds/deletes are *fire-and-forget at the per-document level*: `update.StreamingSolrClients` maintains one `ConcurrentUpdateJettySolrClient` per destination URL (queue size 100, deliberately low thread count — the class comments that more threads "could cause updates to be reordered on a greater scale"), and the document is queued into its stream without waiting for a response. |
| |
| Exceptions that are sent synchronously (blocking per command): in-place updates (a dependent in-place update must not overtake its predecessor in a stream), and forwards to sub-shard leaders or routing-rule targets during shard split / migrate. |
| |
| The leader therefore does not know a document's replica outcome at the time it processes the next document. |
| All outcomes are collected at request end. |
| |
| === Request finish: acknowledgment and error triage |
| |
| When the loader has fed all commands, `DistributedUpdateProcessor.finish()` runs `doDistribFinish()` and then `RunUpdateProcessor.finish()`. |
| `DistributedZkUpdateProcessor.doDistribFinish()` is where the acknowledgment semantics live: |
| |
| 1. If this leader changed its index and skipped any known-stale replicas, bump terms now: `ZkShardTerms.ensureTermsIsHigher(leader, skippedCoreNodeNames)`. |
| 2. `cmdDistrib.finish()` — *block* until every queued replica request has completed (this is the only wait for replica responses in the whole path). |
| 3. Walk `cmdDistrib.getErrors()` and triage: |
| * Error on a `ForwardNode` (this node → leader): added to `errorsForClient`; the client sees the failure and may retry. |
| * Error on a `StdNode` (leader → follower): *not* a client error — the code comments "for now we don't error - we assume if it was added locally, we succeeded". |
| Unless the error is a commit (`commit_end_point` requests never trigger recovery) the follower's coreNodeName is collected for demotion — after double-checking against ZK that we are still the leader and the errored node is still one of our replicas. |
| * Special case: if the remote error's metadata says `cause=LeaderChanged` (SOLR-6511 — the "follower" now believes it is the leader), the error *is* propagated to the client so it can retry against the new leader. |
| 4. `ensureTermsIsHigher(leader, replicasShouldBeInLowerTerms)` — the demotion. |
| The demoted replica's term watcher notices it is behind and puts the core into recovery. |
| This term mechanism (SOLR-11702) is the replacement for the old znode-based "leader-initiated recovery" (LIR); no LIR znodes exist anymore. |
| 5. Compute the achieved replication factor: each shard leader counts itself plus each follower that acked (`LeaderRequestReplicationTracker`), the originating node takes the minimum across shards (`RollupRequestReplicationTracker`), and the result is reported as `rf` in the response header. |
| It is purely informational; a top-of-class TODO ("optionally fail if n replicas are not reached...") records the unimplemented alternative. |
| |
| Client-visible errors are aggregated into `DistributedUpdatesAsyncException` (status: the common code if all agree, else 400 if all 4xx, else 500). |
| |
| Finally `RunUpdateProcessor.finish()` calls `UpdateLog.finish(null)` — the per-request tlog flush described next. |
| |
| === Durability: UpdateLog and TransactionLog |
| |
| `update.UpdateLog` owns a current `update.TransactionLog` (tlog) plus recent old ones, and an in-memory map from doc id to a `LogPtr` into the tlog — the map that makes uncommitted documents visible to realtime get and to atomic-update resolution. |
| |
| Writes: `UpdateLog.add/delete/deleteByQuery` append a record to the current tlog through a buffered `FastOutputStream` over a file channel. |
| *Nothing is flushed per document.* |
| Durability happens per *request*, in `TransactionLog.finish(syncLevel)`: |
| |
| * `NONE` — do nothing. |
| * `FLUSH` (the default) — flush the JVM buffer to the OS; survives a JVM crash / `kill -9`, but not an OS crash or power loss. |
| * `FSYNC` — additionally `channel.force(true)`. |
| The fsync is deliberately outside the buffer lock; the code notes a partial last record after power failure is expected and tolerated by the reader. |
| |
| `syncLevel` is configured on `<updateLog>` in `solrconfig.xml`. |
| So the ref-guide statement "documents are written to the tlog before the indexing call returns" is true, but with default `FLUSH` the response does not imply the bytes reached the disk platter. |
| |
| Commits rotate the tlog: `UpdateLog.preCommit` starts a new tlog (so the old one is definitely fully covered by the index commit), and `postCommit` writes a commit marker into the old one. |
| On startup, `UpdateLog.recoverFromLog()` replays any tlog tail not covered by a commit — this is what makes acked-but-uncommitted updates survive a restart. |
| Retention is bounded by `numRecordsToKeep` (default 100) and `maxNumLogsToKeep` (default 10), which also bound how far a replica can fall behind before PeerSync is impossible and full replication is required. |
| |
| The UpdateLog also has a state machine (`ACTIVE`, `BUFFERING`, `APPLYING_BUFFERED`, `REPLAYING`) used during recovery and shard split: while a core is recovering, incoming `FROMLEADER` updates are written to a separate buffer tlog *without* being applied, and replayed at the end (`applyBufferedUpdates`). |
| State transitions quiesce all in-flight updates through `UpdateLocks.blockUpdates()` (the write side of a fair read/write lock; every normal update holds the read side). |
| See `dev-docs/shard-split/shard-split.adoc` for the shard-split use of buffering. |
| |
| === Replica-side processing |
| |
| A replica receiving `DistribPhase.FROMLEADER` runs the same `versionAdd`/`versionDelete` but with `leaderLogic == false`: |
| |
| * An update without a `\_version_` is rejected (`missing _version_ on update from leader`) — replicas never mint versions. |
| * If the local UpdateLog is not `ACTIVE` (the core is recovering), the update is written to the buffer tlog and dropped (no index write). |
| * Otherwise the *drop rule* runs — the single check that makes asynchronous, possibly-reordered delivery safe: |
| + |
| [source,java] |
| ---- |
| Long lastVersion = vinfo.lookupVersion(cmd.getIndexedId()); |
| if (lastVersion != null && Math.abs(lastVersion) >= versionOnUpdate) { |
| // This update is a repeat, or was reordered. We need to drop this update. |
| return true; |
| } |
| ---- |
| + |
| Application on a replica is therefore idempotent (repeats are dropped) and order-insensitive *per document* (an older version arriving late is dropped). |
| Nothing orders updates across different documents. |
| * On a TLOG replica (not currently leader), the command additionally gets `UpdateCommand.IGNORE_INDEXWRITER`: it is recorded in the tlog but not indexed — the index arrives later by segment replication, and the tlog exists so the replica can replay it if elected leader. |
| |
| Deletes store *negative* versions (hence the `Math.abs`), letting a version lookup distinguish "deleted at version v" from "exists at version v" while still ordering both. |
| |
| Two reorder edge cases get dedicated machinery: |
| |
| * *Delete-by-query*: on the leader, `versionDeleteByQuery` runs under `UpdateLocks.blockUpdates()` — DBQ quiesces *all* updates on the core, because it can affect any document. |
| Replicas keep a list of recent DBQs and re-execute them over an add that arrives out of order relative to the DBQ. |
| A DBQ is also fanned out from the originating node to *all* shard leaders, and is not atomic across shards. |
| * *In-place updates* carry `distrib.inplace.prevversion`; a replica that has not yet seen that previous version waits for it (`waitForDependentUpdates`, using the per-doc lock's `Condition`), and if it never arrives fetches the full document from the leader (`fetchFullUpdateFromLeader`). |
| |
| == Versioning and optimistic concurrency |
| |
| === The _version_ clock |
| |
| `update.VersionInfo.getNewClock()` implements a time-based Lamport clock, synchronized per core: |
| |
| [source,java] |
| ---- |
| long time = System.currentTimeMillis(); |
| long result = time << 20; |
| if (result <= vclock) { |
| result = vclock + 1; |
| } |
| vclock = result; |
| ---- |
| |
| Properties that matter: |
| |
| * Strictly increasing per core, so per-document last-writer-wins is well defined under a single leader. |
| * Wall-clock based so that a restarted or newly elected leader (with an empty in-memory clock) does not go back in time relative to versions already in the index — correctness across leader changes leans on cluster clocks being roughly synchronized. |
| The low 20 bits are a same-millisecond counter (~1M versions/ms before the clock runs ahead of real time). |
| * Not contiguous — a commented-out pure-counter alternative in `VersionInfo` notes contiguous versions would make missing-update detection easier; Solr instead detects gaps via PeerSync's version-list exchange. |
| |
| The `\_version_` field must exist in the schema, single-valued, indexed-or-docValues and stored-or-docValues (`VersionInfo.getAndCheckVersionField`). |
| It must be assigned by Solr internally: user-supplied values would break the replica drop rule. |
| (Use `DocBasedVersionConstraintsProcessorFactory` for application-level version fields.) |
| |
| === Where the version constraint comes from |
| |
| `versionOnUpdate` is taken, in priority order, from the command itself, the document's `\_version_` field, or the `\_version_` request parameter. |
| On internal `FROMLEADER` hops it is the leader-assigned version; on client requests it is the client's optimistic-concurrency constraint (0 when absent). |
| One subtlety: a leader receiving a document forwarded from *another collection* (`distrib.from.collection`, the MIGRATE path) discards the incoming version and stamps its own. |
| |
| === Leader-side optimistic concurrency |
| |
| When a client supplies a nonzero version, the leader checks it against `vinfo.lookupVersion(id)` (tlog first, then index — uncommitted state counts) before assigning the new version: |
| |
| * `> 1` — must equal the current version exactly. |
| * `1` — the document must exist (any positive current version). |
| * `< 0` — the document must not exist. |
| * `0` — no check. |
| |
| A failed check raises `ErrorCode.CONFLICT` (HTTP 409) — unless `failOnVersionConflicts=false`, which silently drops the update instead (useful for batch loads where any conflicting doc should just be skipped). |
| The check runs only on the leader, under the per-document lock, so it is atomic with the version assignment: two clients doing conditional updates on the same document cannot both win. |
| |
| === UpdateLocks |
| |
| `update.UpdateLocks` (SOLR-14679) replaced the historical fixed-size `VersionBucket` striping. |
| It keeps a hash-keyed map of pooled, fair `ReentrantLock`+`Condition` pairs with refcounting, so a lock exists only while some thread is operating on that document id, plus the global `blockUpdatesLock` read/write lock described earlier. |
| Lock acquisition times out after `docLockTimeoutMs`, surfacing pathological contention as an error rather than a hang. |
| |
| Everything that must be atomic per document happens inside `runWithLock`: the OCC check, atomic-update read-modify-write, version assignment, tlog append and index write. |
| This is also what makes realtime get reliable — a concurrent RTG cannot observe a state between the tlog map update and the version assignment. |
| |
| == References |
| |
| * Ref guide: https://github.com/apache/solr/blob/main/solr/solr-ref-guide/modules/deployment-guide/pages/solrcloud-update-consistency.adoc[SolrCloud Update Consistency Model] (the user-facing contract), plus the pages on shards and indexing, recoveries and write tolerance, commits and transaction logs, partial document updates, and realtime get. |
| * `dev-docs/shard-split/shard-split.adoc` — tlog buffering during shard split. |
| * JIRA: https://issues.apache.org/jira/browse/SOLR-11702[SOLR-11702] (shard terms replace LIR), https://issues.apache.org/jira/browse/SOLR-14679[SOLR-14679] (`UpdateLocks` replaces version buckets), https://issues.apache.org/jira/browse/SOLR-6511[SOLR-6511] (`LeaderChanged` propagation), https://issues.apache.org/jira/browse/SOLR-7141[SOLR-7141] (recovery vs. in-flight updates). |
| |
| Not covered here (candidates for future documents): replica-type internals (NRT/TLOG/PULL), recovery and leader election (`RecoveryStrategy`, `PeerSync`, `ZkShardTerms` invariants, `leaderVoteWait`), and commit/visibility internals. |