System Design
X3

Consistency models

A consistency model says which values a read may return when there are many copies and many clients. Stronger models look more like one machine, and cost more.

Not startedSaved in this browser only.
  1. 1Linearizable: one copy in real time. Sequential: one order, real time ignored. Causal: causes before effects.
  2. 2Causal consistency equals the four session guarantees together, and it can stay available in a partition.
  3. 3CAP: during a partition, choose availability or linearizability. PACELC: otherwise, choose latency or consistency.
  4. 4Across regions, a linearizable operation costs at least one round trip to a majority: tens of ms.
X3
    A

    What the user sees

    recorded history
    0102030405060msclient Aclient Bclient CW x=1, in flightR x→1new value, done at 15 msR x→0old value, started at 20 msgap✕ C reads after B finished, yet sees the older valueLinearizable refuses this history. Sequential and causal allow it.
    • A write takes time to reach every copy.
    • A read on a lagging copy returns an older value.
    • Each model draws the line at a different place.

    With copies, two users can see the same data in different states. The consistency model says which states are allowed.

    B

    The models, strongest first

    an arrow: every history above is also allowed below
    Linearizableone copy, real-time order✕ unavailablein a partitionSequentialone order, real time ignoredCausal= all four session guarantees▲ available if theclient stays stickyread yourwrites▲ sticky clientmonotonicreads✓ availablemonotonicwrites✓ availablewritesfollow reads✓ availableEventualcopies agree once writes stop✓ availablestrongerweaker

    Availability in a partition from Bailis et al., Highly Available Transactions, VLDB 2014. Sticky: the client keeps using one copy that has its writes.

    Linearizable is the strongest and cannot stay available in a partition. Causal can, if each client keeps using one copy.

    C

    Which anomalies each model allows

    computed by the checker on the example histories
    modelstale readnew, then oldtwo keysconcurrent ordersown write lostgoes backwrites reorderedreply first
    LinearizablePreventedPreventedPreventedPreventedPreventedPreventedPreventedPrevented
    SequentialAllowedAllowedPreventedPreventedPreventedPreventedPreventedPrevented
    CausalAllowedAllowedAllowedAllowedPreventedPreventedPreventedPrevented
    Read your writesAllowedAllowedAllowedAllowedPreventedAllowedAllowedAllowed
    Monotonic readsAllowedAllowedAllowedAllowedAllowedPreventedAllowedAllowed
    Monotonic writesAllowedAllowedAllowedAllowedAllowedAllowedPreventedAllowed
    Writes follow readsAllowedAllowedAllowedAllowedAllowedAllowedAllowedPrevented
    EventualAllowedAllowedAllowedAllowedAllowedAllowedAllowedAllowed

    Each column is one example history in the checker below. Every model refuses a value that no client wrote.

    I name the anomaly the product cannot show, then pick the weakest model that prevents it.

    D

    Strong or eventual, in product terms

    what the user notices
    user actionwith an eventual readmodel it needsstatus
    Save a profile, then reloadThe old profile comes back.Read your writesSession
    Scroll a feed, then refreshA post appears, then disappears.Monotonic readsSession
    Read a threadA reply shows without its post.CausalCausal
    Buy the last seatTwo people buy it.LinearizableStrong
    Take a unique usernameTwo accounts get the same name.LinearizableStrong
    Revoke access, then test itThe old permission still works.LinearizableStrong
    Count likes on a postThe count is a little behind.EventualStrong not needed

    I use strong reads where a stale answer costs money or trust, and session guarantees everywhere a user only needs to see their own actions.

    E

    Try it: the history checker

    recorded from the lab checker
    history

    A changes a price. B sees the new price and finishes. C reads after that and sees the old price.

    model

    One copy, in real time. Each operation takes effect at one instant between its start and its end.

    01020304050msclient Aclient Bclient CW x=1R x→1R x→0
    ✕ Not allowed. B's R x→1 saw the new value and ended before C's R x→0 started. A later read cannot return the older value.

    Red: the smallest set of operations that breaks the model.

    The checker searches every order for linearizable and sequential, and follows causal chains for the rest. Values start at 0.

    To test a claim, I draw the history. Then I ask if one order of the operations explains every read and keeps the rules of the model.

    F

    Where each read goes

    click a step; its path lights up
    Clientkeeps a tokenService + routerpicks the copyLeaderorders every writeFollowersame regionFolloweracks the writeFollowerother region

    Step 1: Linearizable write

    • Send every write to the leader.
    • The leader waits until a majority of copies has the write, then acknowledges.

    If it fails

    The leader cannot reach a majority: the write blocks or fails. This side of a partition chooses consistency.

    Writes go to the leader and wait for a majority. Each read picks its copy by the guarantee it needs: leader for strong, a follower with a token for session, any follower for eventual.

    G

    CAP, stated precisely

    only during a partition
    side 1side 2client Xclient Yreplica 1x = 2 (new)replica 2x = 1 (old)write x = 2read xlink cut: no messagesChoose A: answerReplica 2 returns x = 1.Y gets an answer, but a stale one.not linearizableChoose C: refuse or waitReplica 2 cannot confirm x is latest.Y gets an error or a timeout.not available
    • C is linearizability: every read returns the latest write.
    • A: every request to a live copy gets an answer. A slow answer still counts.
    • P is a fact of networks, not a choice. The choice is C or A while it lasts.

    Gilbert and Lynch, 2002, prove the theorem for a read/write object with atomic (linearizable) consistency.

    During a partition, a copy that cannot reach the others must either answer with what it has or refuse. Without a partition, I can have both.

    H

    PACELC

    partition: A or C; else: L or C
    Partition now?yes (P)no (E, else)Rare, but forcedA: answer from the local copyC: refuse until the copies meetthe CAP choiceEvery request, every dayL: answer now, maybe staleC: wait for the other copiesa round trip per strong op
    settingpartitionelse
    Cassandra ONE for reads and writesAL
    Cassandra QUORUM, R + W > NCC
    DynamoDB global tables, last writer winsAL
    DynamoDB global tables, multi-Region strongCC
    Postgres one leader, synchronous standbyCC
    Postgres reads on asynchronous standbysAL

    Abadi, 2012: if there is a partition (P), trade availability (A) and consistency (C); else (E), trade latency (L) and consistency (C). Many stores choose per request.

    Most of the time there is no partition, so the real daily trade is latency against consistency.

    I

    Capabilities used

    what each store offers, from its documentation
    toolcapabilitywhat it gives this designalso used for
    PostgresOne primary takes every write; reads on it see each commitLinearizable reads and writes, as long as all traffic goes to the primary.Counters, unique names
    Postgressynchronous_commit = remote_applyThe commit waits until the standby has applied it, so a standby read sees it.Read-your-writes without tokens
    Postgrespg_last_wal_replay_lsn() on a standbyA session token: the standby can tell whether it has applied a given commit.Monotonic reads across standbys
    RedisWAIT numreplicas timeoutLimit The documentation states that WAIT does not make Redis a strongly consistent store.Fewer lost writes on failover
    etcdLinearizable reads by default; serializable reads as an optionStrong reads through consensus, or faster reads that may be stale.Leader election, configuration
    ZooKeeperUpdates from a client apply in the order sent; sync() before a readPer-client order. A read after sync() sees the latest update.Locks, membership
    DynamoDBConsistentRead = trueThe most up-to-date data. Not supported on global secondary indexes.Read-after-write on one table
    DynamoDBGlobal tables: last writer wins, or multi-Region strong (3 Regions)Choose per table: local writes that can conflict, or writes copied to another Region first.Multi-region apps
    CassandraConsistency level per query; lightweight transactions on PaxosR and W per request. IF NOT EXISTS is linearizable, at a higher cost.Unique inserts
    MongoDBCausally consistent sessionsAll four session guarantees, with majority read and write concern.Read-after-write on secondaries
    Cosmos DBFive levels: strong, bounded staleness, session, consistent prefix, eventualSession is the default: read-your-writes for the client that holds the session token.Global apps
    SpannerExternal consistency, using TrueTimeTransactions behave as if run one at a time, in real-time order.Global ledgers
    CockroachDBSerializable by default; no stale readsClose to strict serializability, but not quite: clock skew limits the real-time guarantee.Distributed SQL
    S3Strong read-after-write for PUT and DELETE (since December 2020)A read after a successful write returns that write.Object storage as a source of truth

    I pick the store by the guarantees it documents: a consistent read flag, a consistency level per query, a causal session, or consensus by default.

    J

    Checking linearizable and sequential

    pseudo code
    find an orderpseudo code
    find_order(history, must_precede):     // depth-first search
      IF every op is placed: RETURN order        // the witness1
      FOR EACH op not yet placed:
        IF an unplaced op must_precede it: skip
        IF op is a read AND reg[op.key] != op.value2: skip
        place op; IF op is a write: reg[op.key] = op.value
        IF find_order(rest) succeeds: RETURN it
        undo op                                  // backtrack
      remember (placed set, reg) as a dead end3
      RETURN none
    
    linearizable: must_precede(a, b) = a.end < b.start4
    sequential:   must_precede(a, b) = same client AND a first5
    1. 1The order found is the proof. The checker numbers the operations in this order.
    2. 2A read can come next only if the register holds the value it returned.
    3. 3The same ops placed with the same register values fail the same way. Skip them next time.
    4. 4An operation that ended before another started comes first. Overlapping operations go either way.
    5. 5Only each client’s own order counts. Real time between clients does not.
    Tested source Go: the order search
    Go: the order searchgo
    // findOrder looks for one total order of all ops in which each read returns the latest write on
    // its key (0 if none), and op a comes before op b whenever mustPrecede(a, b). It tries every op
    // whose predecessors are already placed, applies it to the registers, and backtracks on a read
    // that does not match. Dead ends are remembered by (ops placed, register values).
    func findOrder(h History, mustPrecede func(a, b Op) bool) ([]int, bool) {
      n := len(h.Ops)
      placed := make([]bool, n)
      order := make([]int, 0, n)
      regs := map[string]int{}
      dead := map[string]bool{}
    
      var place func() bool
      place = func() bool {
        if len(order) == n {
          return true
        }
        memo := stateKey(placed, regs)
        if dead[memo] {
          return false
        }
        for i, o := range h.Ops {
          if placed[i] || !ready(h, placed, o, mustPrecede) {
            continue
          }
          if o.Kind == Read && regs[o.Key] != o.Value {
            continue // this read cannot come next: the register holds another value
          }
          prev, had := regs[o.Key]
          if o.Kind == Write {
            regs[o.Key] = o.Value
          }
          placed[i] = true
          order = append(order, i)
          if place() {
            return true
          }
          order = order[:len(order)-1]
          placed[i] = false
          if o.Kind == Write {
            if had {
              regs[o.Key] = prev
            } else {
              delete(regs, o.Key)
            }
          }
        }
        dead[memo] = true
        return false
      }
      if !place() {
        return nil, false
      }
      return order, true
    }
    
    // ready reports whether every op that must precede o is already placed.
    func ready(h History, placed []bool, o Op, mustPrecede func(a, b Op) bool) bool {
      for j, p := range h.Ops {
        if !placed[j] && j != o.ID && mustPrecede(p, o) {
          return false
        }
      }
      return true
    }
    
    • The lab compares this search with a check of every permutation on 400 random histories.
    • Each refusal shows the smallest set of operations that has no valid order.

    A history is linearizable if one total order explains every read and keeps real-time order. The order itself is the proof.

    K

    Checking causal

    pseudo code
    causal and sessionspseudo code
    before = each client's own order
           + write -> each read that returns it1
           closed under transitivity2
    
    causal(history):
      FOR EACH read r, FOR EACH write w on r.key:
        IF w before r AND r returned something older than w:
          RETURN refused (w, r)                  // r missed w3
      RETURN allowed
    
    // how the chain reaches r's client names the session guarantee4:
    //   read your writes:     r's client wrote w
    //   monotonic reads:      r's client had read w
    //   monotonic writes:     it read a later write by w's client
    //   writes follow reads:  it read a write made after w was read
    1. 1Reading a value makes the reader depend on the write.
    2. 2Chains count. A reply to a reply depends on the first post.
    3. 3The read returned the initial value, or a write that happens before w.
    4. 4On random histories, the lab finds causal equal to all four guarantees together.
    Tested source Go: causal · Go: session guarantees
    Go: causalgo
    // Causal checks causal consistency: no read returns a value that its own causal past has
    // overwritten. Concurrent writes may appear in any order, and different clients may disagree.
    func Causal(h History) Verdict {
      if r := h.thinAir(); r >= 0 {
        return Verdict{Bad: []int{r}, Rule: RuleThinAir}
      }
      c := newCausalOrder(h)
      for _, a := range h.Ops {
        if c.before[a.ID][a.ID] {
          return Verdict{Bad: []int{a.ID}, Rule: RuleCycle}
        }
      }
      for _, r := range h.Ops {
        if r.Kind != Read {
          continue
        }
        for _, w := range h.Ops {
          // A newer write on the same key happens before the read, and the read returned
          // something older: the initial value, or a write that happens before w.
          if w.Kind != Write || w.Key != r.Key || !c.before[w.ID][r.ID] || !c.older(c.src[r.ID], w.ID) {
            continue
          }
          if c.src[r.ID] < 0 {
            return Verdict{Bad: []int{w.ID, r.ID}, Path: c.path(w.ID, r.ID), Rule: RuleInitRead}
          }
          return Verdict{Bad: []int{c.src[r.ID], w.ID, r.ID}, Path: c.path(w.ID, r.ID), Rule: RuleOverwritten}
        }
      }
      return Verdict{OK: true}
    }
    
    Go: session guaranteesgo
    // ReadYourWrites: after a client writes a key, its reads of that key never return anything older.
    func ReadYourWrites(h History) Verdict {
      if r := h.thinAir(); r >= 0 {
        return Verdict{Bad: []int{r}, Rule: RuleThinAir}
      }
      c := newCausalOrder(h)
      for _, w := range h.Ops {
        for _, r := range h.Ops {
          if w.Kind == Write && r.Kind == Read && r.Key == w.Key && programBefore(w, r) && c.older(c.src[r.ID], w.ID) {
            return Verdict{Bad: []int{w.ID, r.ID}, Path: []int{w.ID, r.ID}, Rule: RuleRYW}
          }
        }
      }
      return Verdict{OK: true}
    }
    
    // MonotonicReads: once a client has read a value, its later reads of that key never return
    // anything older.
    func MonotonicReads(h History) Verdict {
      if r := h.thinAir(); r >= 0 {
        return Verdict{Bad: []int{r}, Rule: RuleThinAir}
      }
      c := newCausalOrder(h)
      for _, r1 := range h.Ops {
        for _, r2 := range h.Ops {
          if r1.Kind == Read && r2.Kind == Read && r1.Key == r2.Key && programBefore(r1, r2) &&
            c.src[r1.ID] >= 0 && c.older(c.src[r2.ID], c.src[r1.ID]) {
            return Verdict{Bad: []int{r1.ID, r2.ID}, Path: []int{c.src[r1.ID], r1.ID, r2.ID}, Rule: RuleMR}
          }
        }
      }
      return Verdict{OK: true}
    }
    
    // MonotonicWrites: a client that sees a client's second write also sees its first.
    func MonotonicWrites(h History) Verdict {
      if r := h.thinAir(); r >= 0 {
        return Verdict{Bad: []int{r}, Rule: RuleThinAir}
      }
      c := newCausalOrder(h)
      for _, w1 := range h.Ops {
        for _, w2 := range h.Ops {
          if w1.Kind != Write || w2.Kind != Write || !programBefore(w1, w2) {
            continue
          }
          if v, ok := sawThenMissed(h, c, w2.ID, w1); ok {
            return Verdict{Bad: append([]int{w1.ID}, v...), Path: append([]int{w1.ID}, v...), Rule: RuleMW}
          }
        }
      }
      return Verdict{OK: true}
    }
    
    // WritesFollowReads: a client's write comes after everything the client had read, directly or
    // through earlier reads. A client that sees the write also sees those values.
    func WritesFollowReads(h History) Verdict {
      if r := h.thinAir(); r >= 0 {
        return Verdict{Bad: []int{r}, Rule: RuleThinAir}
      }
      c := newCausalOrder(h)
      for _, w1 := range h.Ops {
        for _, w2 := range h.Ops {
          // w1 is another client's write that w2's writer had seen before writing w2.
          if w1.Kind != Write || w2.Kind != Write || w1.Proc == w2.Proc || !c.before[w1.ID][w2.ID] {
            continue
          }
          if v, ok := sawThenMissed(h, c, w2.ID, w1); ok {
            p := append(c.path(w1.ID, w2.ID), v[1:]...)
            return Verdict{Bad: append([]int{w1.ID}, v...), Path: p, Rule: RuleWFR}
          }
        }
      }
      return Verdict{OK: true}
    }
    
    // sawThenMissed finds a client that reads write w2 and later reads w1's key and gets something
    // older than w1. It returns [w2, that read, the later read].
    func sawThenMissed(h History, c causalOrder, w2 int, w1 Op) ([]int, bool) {
      for _, r := range h.Ops {
        if r.Kind != Read || c.src[r.ID] != w2 {
          continue
        }
        for _, r2 := range h.Ops {
          if r2.Kind == Read && r2.Key == w1.Key && programBefore(r, r2) && c.older(c.src[r2.ID], w1.ID) {
            return []int{w2, r.ID, r2.ID}, true
          }
        }
      }
      return nil, false
    }
    

    Causal consistency forbids one thing: a read that returns a value its own causal past has overwritten.

    L

    Failure cases

    what breaks, and what saves you
    eventresultwhy it stays correctsaved by
    A user reads a follower right after a writeThe follower may not have it.The token makes the follower wait, or the read goes to the leader.Session token
    The load balancer sends the next read to another followerThat follower can be further behind.The token carries the newest position seen, so no read goes back.Session token
    A partition cuts the leader from the majorityThe leader cannot commit.Writes on the minority side fail. The majority side elects a leader and goes on.Quorum
    An old leader wakes after a pauseIt may serve reads it no longer owns.It checks with a majority, or its lease expired before the new leader started.Lease, quorum read
    Two regions write the same key at onceA conflict.Give each key a home region, or merge with a CRDT. Last writer wins drops one write.Home region
    A cache sits in front of a strong storeThe cache returns old values.Reads that must be strong skip the cache. Writes invalidate the key.Bypass
    A reply is replicated before its postA reader sees the reply alone.The reply carries the post's position. A copy shows it only after it has the post.Dependency token
    M

    Scale ladder

    start simple; climb only on a signal
    Each step adds one component1One leader2+ followers, tokens3Leaderless quorum4Multi-regionmore load →
    Latency of one consistent operation0.11101001kCommit, one leader: 0.14 msCommit, one leader0.14+ sync standby, local: 0.36 ms+ sync standby, local0.36Majority, Virginia + Oregon: 35 msMajority, Virginia + Oregon35Dublin to Virginia leader: 55 msDublin to Virginia leader55Mumbai to Virginia leader: 129 msMumbai to Virginia leader129ms, log scale
    stepaddit handlesmove up when you see
    1One leader for all reads and writes, with a standby for failover.Linearizable for free. About 107,000 primary-key reads a second on one Postgres in the lab.Reads load the leader.
    2Followers for reads, with session tokens. Strong reads stay on the leader.Reads grow with each follower. Users see their own writes and never go back.A failover pause is not acceptable, or writes come from many places.
    3A leaderless quorum: N copies, R + W > N for strong reads.No failover step. Each request picks its own R and W.Users on several continents need local latency.
    4Several regions: local reads with session or causal guarantees.Region loss. Only strong operations pay a round trip to a majority: 35 ms or more.Top of the ladder. Shard so each key has a home region.

    Local latencies measured in the lab for the replication sheet. Regional figures are floors: great-circle distance at about 200 km per ms in fibre, there and back. Real cables run longer.

    I start with one leader that serves everything, so every read is linearizable. I add followers with session tokens for reads, and pay cross-region round trips only for operations that need them.

    N

    Drill

    predict, then reveal

    0 of 10 known

    1. A write of x = 1 ends. Later, another client reads x = 0. Is the history linearizable? Sequential?

    2. One user sees a new price. A second user, who loads the page after the first finished, sees the old price. Which model forbids this?

    3. A user saves a change and reloads. The change is missing. Which guarantee broke, and what are two fixes?

    4. A reply appears before the post it answers. Which guarantee broke?

    5. State CAP precisely.

    6. There is no partition. Why does a multi-region store still trade consistency?

    7. Can a store stay causally consistent during a partition?

    8. Why is it useful that linearizability is local?

    9. Are serializable and linearizable the same?

    10. An old leader pauses, a new one is elected, and the old one wakes. How do you keep its reads linearizable?

    O

    Numbers to say

    derived or cited
    fibre
    Light covers about 200,000 km a second in fibre: 200 km per ms. Refractive index about 1.47.
    US coasts
    Virginia to Oregon is 3,500 km: 17.5 ms one way, 35 ms there and back.
    Atlantic
    Virginia to Dublin is 5,500 km: 27 ms one way, 55 ms there and back.
    Asia
    Virginia to Mumbai is 12,900 km: 64 ms one way, 129 ms there and back.
    majority
    A strong write waits for the nearest majority: one round trip to the closest other replica, at least.
    local
    0.14 ms per commit on one Postgres; 0.36 ms with a synchronous standby on the same machine.

    Distances are great-circle; round trips are floors. Local commits measured in the lab on an 8-core laptop.