Skip to content
KEDBYTE
Site navigation
How Data Works
Chapter
36

Consistency, Availability and Network Partitions

Part F · Data Across Machines|3,753 words|about 16 min read|Volume F

36.0 What this chapter gives you#

  1. “The system is consistent” is incomplete unless we state which observations are allowed. A database can obey its local constraints while different clients observe incompatible versions of a shared value.
  2. You will learn the meanings needed to reason about a network split: linearizability, weaker read contracts, availability, quorum intersection and the limits captured by CAP.
  3. The examples are small histories and set calculations. They are not rankings of database products and do not prove that a real distributed deployment satisfies a consistency model.

36.1 Name the consistency model#

36.1.1 PLAIN — in simple words#

  1. A consistency model describes what different readers and writers may observe. It is a contract about histories of operations, not simply whether all copies eventually contain the same bytes.
  2. A strong single-object contract can make the object appear to change at one instant between each request and its reply. If one write finishes before a later read starts, that read cannot pretend the completed write never happened unless a later valid change explains the result.
  3. Weaker contracts may allow older values or different observations under stated rules. They can be useful, but only when the application knows what it is accepting.

36.1.2 PLAIN — a picture in your head#

  1. Imagine a single reservation board watched by several people. In the strongest simple picture, every operation has one place in a shared timeline that respects actions already completed before another begins.
  2. A system of copied boards may instead promise that everyone will eventually catch up, without promising that each immediate read is current.
  3. Where the comparison breaks: the contract concerns observable operation histories. It does not require a physically instantaneous update to every machine or a universal perfectly synchronized wall clock.

36.1.3 PLAIN — a worked example#

  1. A register starts at 0. Client X writes 1 and receives success. Only afterwards, client Y invokes a read. No other write occurs.
  2. A linearizable read must return 1. Returning 0 cannot be placed after the completed write in a valid single-register history.
  3. If the read overlaps the write instead, either 0 or 1 can be consistent with a suitable ordering, depending on the rest of the history. Overlap is not the same as “the write finished first.”
  4. A later valid write of 0 changes the analysis again. Correctness is judged against the whole relevant history, not a rule that values must always increase.

36.1.4 PLAIN — what is really happening inside#

  1. Real-time precedence means the order of non-overlapping completed operations, not comparing arbitrary timestamps from different machines. The system can reason about messages and completion boundaries without assuming identical clocks.
  2. Sequential consistency requires a shared order respecting each participant’s program order but does not impose every cross-client real-time precedence constraint. Causal consistency preserves relevant cause-and-effect dependencies while allowing more freedom for unrelated work.
  3. Eventual convergence is weaker still unless supplemented with precise guarantees. It needs conditions such as delivery, conflict resolution and an opportunity to finish propagating updates; the word “eventual” is not a deadline.

36.1.5 TECHNICAL — the engineer’s version#

  1. The CAP discussion here uses atomic read/write object semantics, commonly described as linearizability: each operation appears at a point within its invocation-response interval. [S153]
  2. Do not conflate single-object linearizability with transaction serializability or with the C in ACID. Transaction isolation concerns relationships among transactional operations; application consistency concerns invariants. [S66] [S117]
  3. A specification should name object scope, session guarantees, visibility rules, failure behaviour and any staleness bound. “Strong” or “eventual” without those details is not sufficient for a design review.

36.1.6 WORDS — remember these#

  1. Linearizability: operations appear to take effect in one real-time-respecting order — a correctness condition placing each operation within its invocation-response interval. Sequential consistency: one shared order respects each client’s own order — a model that does not additionally require all cross-client real-time precedence. Causal consistency: effects follow their relevant causes — a model preserving causal dependencies while allowing flexibility for concurrent unrelated operations.

36.2 What availability promises#

36.2.1 PLAIN — in simple words#

  1. Availability in an operational dashboard often means the proportion of requests answered acceptably within a deadline. CAP uses a theoretical requirement: requests at non-failed participants eventually receive the required responses.
  2. These are different measurements. A response after an hour might satisfy an eventual-response property while being useless for a checkout that must finish in two seconds.
  3. Returning an error to every valid request does not magically satisfy the intended read/write service. The response must be interpreted under the service specification rather than counted as success merely because some bytes arrived.

36.2.2 PLAIN — a picture in your head#

  1. A shop that promises to answer every customer eventually has made a different promise from one that promises to answer 99.9% within two seconds.
  2. A sign saying “unavailable” is a response in everyday language, but it does not provide the requested reservation or current stock answer.
  3. Where the comparison breaks: formal availability is a property quantified over executions and requests. A month’s observed uptime is empirical evidence for an operational target, not a proof over all possible executions.

36.2.3 PLAIN — a worked example#

  1. Consider a hypothetical one-hour test with 10,000 eligible reads. Of these, 9,990 return valid results within the agreed deadline, five return late and five fail.
  2. Deadline-qualified success is 9,990 / 10,000 = 99.9%. Counting the late responses as acceptable would produce a different metric and must be declared.
  3. That finite test does not establish CAP availability under every possible partition. It records the behaviour of one workload during one observation window.
  4. Separate the denominator, deadline, excluded requests and response-validity rule before comparing systems or reporting service health.

36.2.4 PLAIN — what is really happening inside#

  1. An operation may require coordination with a remote participant. If the remote path is unavailable, the implementation must wait, use a supported weaker mode or reject according to its contract.
  2. The ability to serve some operations does not imply the ability to serve all operations. A catalogue page may remain available while a stock reservation is deliberately paused.
  3. Availability is therefore often operation-specific. A system-level percentage can hide that the most important action is the one failing.

36.2.5 TECHNICAL — the engineer’s version#

  1. The formal CAP availability condition is eventual completion of requests under the specified service and failure model, not a latency percentile or a commercial uptime percentage. [S153]
  2. Operational service-level indicators should count eligible events under a defined success predicate and time boundary. Their measurement does not establish a universal distributed-systems theorem. [S103]
  3. Document degradation by operation: authoritative writes, stale-tolerant reads, status lookups and administrative changes may have different coordination requirements.

36.2.6 WORDS — remember these#

  1. Theoretical availability: required requests eventually complete — a liveness property under a stated distributed-system model. Service-level indicator: a measured service outcome — a defined quantity such as the fraction of eligible requests meeting a correctness and latency condition. Degraded mode: a documented reduced service — an operational state with explicit limits rather than an undisclosed weakening of guarantees.

36.3 Partitioned communication#

36.3.1 PLAIN — in simple words#

  1. A network partition prevents some participants from communicating with others. The machines on both sides may still be running and receiving client requests.
  2. Each side can lack information about actions accepted on the other. Waiting longer may help when the interruption ends, but no finite wait can guarantee delivery during an indefinitely lasting split.
  3. A system must decide which actions it can safely perform with the information it has. It cannot obtain missing knowledge merely by adding a timeout.

36.3.2 PLAIN — a picture in your head#

  1. Two branches lose their connecting phone line while both remain open. Each still has the last shared stock figure, but neither can learn about new promises made by the other.
  2. If both sell from the same undivided last item, a later conversation may reveal that two customers were promised one object.
  3. Where the comparison breaks: some applications can pre-allocate disjoint rights or use operations designed to merge safely. Those approaches change the coordination requirement rather than making communication failure disappear.

36.3.3 PLAIN — a worked example#

  1. Three replicas A, B and C begin with value 0. The network splits into {A} and {B,C}.
  2. If B and C use an established majority protocol, that side may remain eligible to accept coordinated writes. A cannot obtain a majority by counting itself twice or inventing a replacement member.
  3. A can still answer an explicitly stale-tolerant question from its local copy, but it cannot represent that answer as satisfying a current linearizable read without the necessary evidence.
  4. Membership, quorum and read rules must already define the behaviour. The set diagram alone does not implement them.

36.3.4 PLAIN — what is really happening inside#

  1. Network problems can be asymmetric: A may reach B while B’s replies do not reach A, or one client may reach an old leader while another reaches its replacement.
  2. The system’s authority and consistency mechanisms must handle delayed messages after recovery. Reconnection does not automatically make incompatible accepted histories mergeable.
  3. Avoid changing membership independently on both sides merely to regain a local majority. That can create two disjoint groups each calling itself authoritative.

36.3.5 TECHNICAL — the engineer’s version#

  1. Treat partition tolerance as a failure-model requirement: communication may be delayed or lost between components. It is not a third optional feature that can be removed from a real network by a configuration slogan. [S153]
  2. Consensus membership changes need rules that preserve the relevant quorum intersection. The Raft paper’s joint-consensus design requires majorities of both configurations during its transition. [S154]
  3. Failure-injection evidence should specify direction, duration, affected paths, clock assumptions and recovery behaviour. Killing a process is not the same experiment as partitioning live participants.

36.3.6 WORDS — remember these#

  1. Network partition: some participants cannot communicate — a communication failure that can separate live components. Asymmetric failure: reachability differs by direction or path — a condition in which one observer’s connectivity does not describe the whole system. Membership configuration: the participants counted by the protocol — a versioned set whose changes require safety-preserving rules.

36.4 CAP with precise definitions#

36.4.1 PLAIN — in simple words#

  1. CAP identifies a limit: under partitioned communication, a shared read/write service cannot guarantee both the specified single-copy current behaviour and eventual completion of every required request at non-failed participants.
  2. The conflict comes from missing information. One side cannot always know whether a completed write occurred on the other side while the two cannot communicate.
  3. This is not a instruction to choose two attractive letters on a product brochure. It is a reason to specify what the application does when coordination is unavailable.

36.4.2 PLAIN — a picture in your head#

  1. Dev must answer a question about a board in a room he cannot contact. In one possible history, Mira has changed the board to 1. In another, she changed it to 2.
  2. Dev sees the same evidence in both histories. A definite answer cannot be guaranteed correct for both, and waiting forever fails the requirement to eventually answer.
  3. Where the comparison breaks: the theorem has precise object and execution definitions. An analogy cannot establish every possible trade-off for every kind of service.

36.4.3 PLAIN — a worked example#

  1. Let replicas A and B be separated. Consider history H1: A completes a write of 1. Consider H2: A completes a write of 2. B receives no messages distinguishing H1 from H2.
  2. A later read reaches B. To satisfy linearizable register semantics, its answer must be 1 in H1 and 2 in H2, assuming no intervening writes.
  3. But B has the same local evidence in both histories. A deterministic choice cannot be correct in both; waiting for information that may never arrive violates eventual-response availability.
  4. This original worked trace illustrates the information conflict discussed by Gilbert and Lynch. It does not claim that every useful distributed service requires this exact register contract. [S153]

36.4.4 PLAIN — what is really happening inside#

  1. A service can preserve the stronger contract by refusing or delaying operations that cannot safely coordinate. It can instead provide explicitly weaker observations or accept operations under a conflict-resolution policy.
  2. Different actions may make different choices. Reading an older product description is not equivalent to promising the same last physical unit twice.
  3. Once the partition ends, reconciliation still needs an application meaning. Convergence to one value does not undo commitments already made to customers.

36.4.5 TECHNICAL — the engineer’s version#

  1. In this chapter C means linearizable or atomic read/write semantics, not schema constraints or ACID consistency. A means the formal availability condition, not a measured success percentage. P describes unreliable communication including partitions. [S153]
  2. The impossibility concerns guarantees across the allowed executions. It does not say that a healthy connected deployment must always choose between correct responses and availability on each ordinary request.
  3. A system may also fail to provide either desired property because of bugs, overload or incomplete design. The theorem is a limit, not a certificate that any chosen implementation provides the remaining letters.

36.4.6 WORDS — remember these#

  1. CAP theorem: a limit on joint guarantees under unreliable communication — the incompatibility of the specified atomic service and availability requirements in the partition-prone model. Indistinguishable histories: different realities look the same to an observer — executions with identical local evidence that require different correct outcomes. Consistency downgrade: a weaker observation contract is used — a change that must be explicit and acceptable for the operation.

36.5 Quorum reasoning and limits#

36.5.1 PLAIN — in simple words#

  1. A quorum is a sufficiently large set of participants for a particular protocol step. Sets can be chosen so that any two relevant quorums overlap.
  2. Overlap means they share at least one participant. It does not, by itself, explain which value that participant stores, which version a reader chooses or how concurrent writes are ordered.
  3. Quorum arithmetic is therefore one ingredient of a protocol, not a complete substitute for one.

36.5.2 PLAIN — a picture in your head#

  1. In a group of five clerks, any two groups of three share at least one clerk. That shared person can connect the groups’ information only if the rules require them to retain and respect the relevant evidence.
  2. If the clerk forgets an accepted decision, gives two incompatible approvals or answers with an old note, the head count alone does not save the procedure.
  3. Where the comparison breaks: real protocols specify durable votes, versions, message handling and membership changes. People standing in overlapping circles are only the set-theoretic part.

36.5.3 PLAIN — a worked example#

  1. With N=3, a write set of size W=2 and a read set of size R=2 must overlap because R+W=4 > 3. The minimum overlap is R+W-N=1 participant.
  2. With N=5, two majorities of size 3 also overlap in at least one participant. The lab enumerates such sets and checks the intersection property directly.
  3. Now weaken the read to R=1 while keeping W=2 among three replicas. A read from the third untouched replica can miss a completed write.
  4. Even R=2,W=2 does not alone prove linearizability. Consider concurrent versions, incomplete writes, repair rules, the selection of the latest value and whether all operations used the same membership.

36.5.4 PLAIN — what is really happening inside#

  1. Protocols attach meaning to the intersection. A voter may remember an accepted term and log, or a storage replica may keep versioned values and participate in a defined read-repair procedure.
  2. “Sloppy” quorum schemes can use substitute nodes outside the usual replica set during failures. The ordinary fixed-set intersection argument cannot simply be copied without examining that change.
  3. Wall-clock timestamps are not automatically a safe ordering mechanism. Clock skew and concurrent updates can make the largest displayed timestamp a poor representation of causal or real-time order.

36.5.5 TECHNICAL — the engineer’s version#

  1. For fixed subsets of an N-member universe, |R ∩ W| >= |R| + |W| - N. The bound proves intersection only; it does not prove an atomic-register implementation.
  2. The original Dynamo design uses version handling, reconciliation and sloppy quorums to pursue its stated availability goals. It should not be summarized as “R+W>N therefore linearizable.” [S162]
  3. Review membership, write completion, version selection, partial failures, concurrent writes and read repair before asserting a consistency guarantee. The book’s set-enumeration tests validate arithmetic, not a complete distributed protocol.

36.5.6 WORDS — remember these#

  1. Quorum: a participant set sufficient for a protocol action — a set meeting defined size or intersection requirements. Quorum intersection: relevant sets share participants — a necessary ingredient in many protocols whose meaning depends on retained state and rules. Sloppy quorum: substitute reachable nodes may count — a failure-handling scheme that differs from always using one fixed replica set.

36.6 Application-level consequences#

36.6.1 PLAIN — in simple words#

  1. The useful question is not “Which consistency word sounds strongest?” It is “Which incorrect or delayed outcomes can this operation tolerate, and which must be prevented?”
  2. The answer may differ between browsing a catalogue, editing a draft, reserving scarce stock and calculating a historical report.
  3. A documented refusal can be safer than a confident but unsupported success. A clearly labelled older answer can also be useful when the task genuinely permits it.

36.6.2 PLAIN — a picture in your head#

  1. A shop can show yesterday’s opening-hours brochure while its branch link is down, provided the reader is told its date. It should not call that brochure a confirmed reservation for the final item.
  2. The same communication outage therefore leads to different behaviour for different tasks.
  3. Where the comparison breaks: the acceptable policy belongs to the application owner and users. A technical mechanism does not decide the commercial or safety consequences on their behalf.

36.6.3 PLAIN — a worked example#

  1. Define three fictional contracts: catalogue descriptions may be up to ten minutes old; a user’s draft read must include their last acknowledged edit; a reservation may succeed only through the current authorized stock owner.
  2. During a split, a cached catalogue may still be eligible if its age and scope meet the contract. A draft read may need to wait or route elsewhere. A reservation on an isolated stale node must not pretend to succeed authoritatively.
  3. The ten-minute allowance is a chosen example policy, not a universal freshness recommendation. Each decision should be tested against the actual harm of stale or unavailable information.
  4. Record the mode in the response and operational evidence so later reports can distinguish authoritative outcomes from permitted stale views.

36.6.4 PLAIN — what is really happening inside#

  1. Consistency requirements propagate through caches, replicas, sessions, transactions and downstream services. A strong database underneath a stale cache does not automatically provide a strong application read.
  2. Tests should include ambiguous boundaries: write completed but reply lost, follower behind, ownership changed, cache invalidation delayed and only part of a report available.
  3. Review the whole request path, including what the user sees. Correct internal refusal followed by a misleading “confirmed” screen is still an application defect.

36.6.5 TECHNICAL — the engineer’s version#

  1. Express each operation’s contract as allowed histories and observable outcomes. Include degraded-mode semantics, retries, deadlines and source-of-authority checks.
  2. Combine model reasoning with actual implementation tests. Neither a CAP explanation nor a passing quorum arithmetic test certifies a production application. [S153] [S154]
  3. Preserve the difference between a proven impossibility under a formal model, a documented product guarantee and an empirical observation from a bounded experiment.

36.6.6 WORDS — remember these#

  1. Operation contract: the promise made for one action — a specification of valid results, timing, visibility and failure outcomes. Authoritative outcome: a result accepted under the valid authority path — a state distinct from a cached guess or tentative local proposal. Allowed history: a sequence the contract permits — an operation trace used to evaluate a consistency or correctness claim.

36.97 Practice and worked answers#

  1. A write of 1 completes before a read begins, with no later writes. Is a returned 0 linearizable? No. It cannot fit the required real-time-respecting register history.
  2. Does linearizability require every machine’s clock to show the same time? No. The property concerns invocation-response ordering, not identical wall-clock displays.
  3. Does 99.9% observed deadline success prove formal CAP availability? No. It describes a finite operational measurement under a defined denominator.
  4. Why cannot an isolated reader distinguish H1’s completed write of 1 from H2’s completed write of 2? The relevant messages are absent in both local observations.
  5. What is the minimum overlap of sets of size 3 and 3 in a universe of 5? At least one member, since 3+3-5=1.
  6. Does that overlap prove linearizability? No. Version, completion, concurrency, failure and membership rules remain necessary.
  7. May both sides independently replace membership until each has a majority? Not safely without a protocol preserving the required cross-configuration intersection.
  8. Why can catalogue reads and reservations have different partition behaviour? They have different correctness and freshness requirements. The distinction must be explicit, not hidden in a generic availability label.

36.98 Common wrong ideas#

  1. Wrong: consistency has one universal meaning. Right: name the observation model and scope.
  2. Wrong: ACID consistency is CAP consistency. Right: these terms address different properties.
  3. Wrong: eventual means within a fixed short delay. Right: a deadline needs an additional explicit bound and assumptions.
  4. Wrong: CAP is always “pick any two features.” Right: it concerns joint guarantees under unreliable communication and precise definitions.
  5. Wrong: every error response counts as useful availability. Right: responses must satisfy the intended service contract.
  6. Wrong: quorum intersection alone proves the newest value is returned. Right: protocol rules give the intersection its meaning.
  7. Wrong: reconnecting automatically repairs business promises. Right: converged storage may still require application reconciliation.
  8. Wrong: a strong database guarantees a strong application. Right: every layer on the request path must preserve the contract.

36.99 Chapter summary in 20 lines#

  1. Consistency models define allowed observations across operations.
  2. Linearizability respects the real-time order of non-overlapping operations.
  3. Overlapping operations may admit more than one valid order.
  4. Sequential consistency and causal consistency impose different constraints.
  5. Eventual convergence is not a fixed freshness deadline.
  6. Transaction isolation and application invariants are separate concepts.
  7. CAP availability is a theoretical eventual-completion property.
  8. Operational availability requires a defined success predicate and denominator.
  9. A network partition can separate participants that remain alive.
  10. Timeouts cannot manufacture missing information.
  11. Membership changes must preserve protocol safety.
  12. CAP exposes an information conflict under precise assumptions.
  13. It is not a product ranking or a certificate of implementation quality.
  14. Fixed-set quorum arithmetic can prove overlap.
  15. Overlap alone does not prove linearizable behaviour.
  16. Version selection, partial writes and membership matter.
  17. Sloppy quorums change the ordinary fixed-replica analysis.
  18. Different operations can justify different degraded modes.
  19. Caches and routing must preserve the intended application contract.
  20. Distinguish formal reasoning, documented guarantees and measured evidence.

Return to contents