System Design
X4

Time and order

Each machine has its own clock, and no two agree. To order events across machines, use the messages between them, or bound the error and wait it out.

Not startedSaved in this browser only.
  1. 1Wall clocks disagree by milliseconds or more. Never order causally related events by wall time.
  2. 2Measure durations with the monotonic clock; use the wall clock only to tell the time.
  3. 3Lamport clocks keep causal order; vector clocks also detect concurrency; HLC stays near wall time.
  4. 4Time-sorted IDs must wait or fail when the clock steps back, or they repeat.
X4
    A

    Last write wins, with skew

    recorded from the lab simulation
    true time →node A+10 msnode B-10 msa1write a1wall 114sees a1write b1wall 99b1 happens after a1. Its wall time, 99, is smaller than 114.✕ Last write wins keeps a1. B's newer write is lost.Skew 20 ms between A and B; the message took 2 ms.
    • Each machine sets its clock over the network, with an error.
    • Skew: two clocks read different times at the same instant.
    • A message can then arrive "before" it was sent.

    A wall clock can be behind, so a later write can get a smaller timestamp. Last write wins then drops the newer value.

    B

    Which clock, for what

    what each stamp can tell you
    clockkeeps causal orderfinds concurrencynear real time
    Wall clockNoNoWithin the skew
    Monotonic clockOne machineNoArbitrary start
    LamportYesNoNo
    VectorYesYesNo
    Hybrid logical (HLC)YesNoWithin the skew
    TrueTime intervalWith commit waitNoWithin ε

    A vector holds one counter per node. A Lamport or HLC stamp fits in 64 bits.

    I use wall time to tell users the time, the monotonic clock for durations, and logical or hybrid clocks to order events across machines.

    C

    Try it: three nodes, skewed clocks

    recorded from the seeded lab simulation
    skew between fastest and slowest clock
    stamp with

    Each event takes the node’s wall clock reading, in ms. Last write wins keeps the write with the largest stamp. Offsets: A 10 ms, B -10 ms, C 0 ms.

    event 12 of 12: A gets m4, true time 132 ms
    events in true-time order →node A+10 msnode B-10 msnode C0 msa1114c1104send m1115recv m197b199send m3110recv m3123send m2103recv m2119c2123send m4127recv m4142
    order by wall clock
    1. B: recv m1 97
    2. B: b1 99
    3. B: send m2 103
    4. C: c1 104
    5. C: send m3 110
    6. A: a1 114
    7. A: send m1 115
    8. C: recv m2 119
    9. A: recv m3 123
    10. C: c2 123
    11. C: send m4 127
    12. A: recv m4 142
    m1 arrives with a smaller stamp than its send. Lost by last write wins: b1 happens after a1 but has the smaller stamp. Last write wins keeps c2; the causally last write is c2.

    Red chips: an event listed before something that happened before it. Seed 1; the same seed gives the same run at every skew.

    With skew, wall-clock order contradicts causality. Lamport and vector clocks ignore the wall clock, so their order is the same at every skew, and vectors also show which writes are concurrent.

    D

    How far apart are the clocks

    cited, as orders of magnitude
    Clock error by sync method101001k10k100kNTP, public internet: tens of msNTP, public internettens of msTrueTime ε, average: 4 ms (1 to 7)TrueTime ε, average4 ms (1 to 7)NTP, fast LAN: a few hundred µsNTP, fast LANa few hundred µsPTP hardware clock: under 40 µsPTP hardware clockunder 40 µsmicroseconds, log scale
    drift
    A clock that is not corrected drifts. TrueTime applies a drift rate of 200 µs per second between syncs.
    sawtooth
    30 s between polls × 200 µs per second = 6 ms. So ε grows from about 1 ms to 7 ms, then drops at the next sync.
    skew
    Two clocks, each within ε of true time, can be up to 2ε apart.

    NTP: RFC 5905 and the NTP project. TrueTime: the Spanner paper, OSDI 2012. PTP: the cloud provider's typical error bound on supported instances.

    NTP over the internet keeps clocks within tens of milliseconds. On a good LAN, a few hundred microseconds. PTP hardware clocks get to microseconds.

    E

    Wall clock and monotonic clock

    illustration
    true time →readsmonotonicNTP steps backwallstartendWall: end − start < 0. Monotonic: always ≥ 0.
    useclockstatus
    Show the time, stamp a log linewallApproved
    Timeout, lease, latencymonotonicApproved
    Timeout from two wall readingswallNot approved
    Order events across machineswallNot approved
    Compare monotonic readings of two machinesmonotonicDifferent starts

    A timeout, a lease or a latency is a duration, so I measure it with the monotonic clock. The wall clock can jump.

    F

    Lamport and vector clocks

    pseudo code
    logical clockspseudo code
    on a local event or a send, at node n:
      lamport[n] += 1
      vector[n][n] += 1                    // count your own event
      stamp the event, and the message, with both
    
    on receive of m, at node n:
      lamport[n] = max(lamport[n], m.lamport) + 11
      vector[n] = entrywise max2(vector[n], m.vector)
      vector[n][n] += 1
    
    a before b   IF a.vector <= b.vector in every entry, and they differ3
    concurrent   IF neither is before the other4
    1. 1The receive gets a larger stamp than the send. Lamport, 1978: if a happens before b, C(a) < C(b).
    2. 2The node now knows everything the sender knew.
    3. 3The lab checks this rule against the true causal order for every pair, over 50 seeds and 4 skews.
    4. 4Concurrent: a conflict for last write wins to hide, or for the app to merge.
    Tested source Go: all four clocks in one run
    Go: all four clocks in one rungo
    // Run plays the script under a timing and a skew. Events come out in true-time order; ties go
    // to the lower node.
    func Run(tm Timing, skew int) ([]Event, error) {
      off := Offsets(skew)
      var (
        pos     [Nodes]int        // next step of each node
        ready   [Nodes]int        // true time at which each node may act next
        lamport [Nodes]int        // each node's Lamport counter
        vector  [Nodes][Nodes]int // each node's vector clock
        hlc     [Nodes]HLC        // each node's hybrid logical clock
        sent    = map[int]msgStamp{}
        events  []Event
      )
      for n := range Nodes {
        ready[n] = Start
      }
      for {
        // Pick the node whose next step can happen first. A recv cannot happen before its
        // message arrives.
        best, bestAt := -1, 0
        for n := range Nodes {
          if pos[n] >= len(script[n]) {
            continue
          }
          s := script[n][pos[n]]
          at := ready[n] + tm.Gap[n][pos[n]]
          if s.kind == KindRecv {
            m, ok := sent[s.msg]
            if !ok {
              continue // not sent yet
            }
            at = max(at, m.real+tm.Delay[s.msg])
          }
          if best < 0 || at < bestAt {
            best, bestAt = n, at
          }
        }
        if best < 0 {
          break
        }
        n, s := best, script[best][pos[best]]
        e := Event{ID: len(events), Node: n, Kind: s.kind, Value: s.value, Msg: s.msg, Peer: -1, Real: bestAt, Wall: bestAt + off[n]}
        if s.kind != KindWrite {
          e.Peer = s.peer
        }
        pt := e.Wall
    
        if s.kind == KindRecv {
          m := sent[s.msg]
          lamport[n] = max(lamport[n], m.lamport) + 1 // Lamport: jump past the sender
          for i := range Nodes {                      // vector: element-wise max
            vector[n][i] = max(vector[n][i], m.vector[i])
          }
          hlc[n] = hlcRecv(hlc[n], m.hlc, pt)
        } else {
          lamport[n]++
          hlc[n] = hlcLocal(hlc[n], pt)
        }
        vector[n][n]++
    
        e.Lamport, e.Vector, e.HLC = lamport[n], vector[n], hlc[n]
        if s.kind == KindSend {
          sent[s.msg] = msgStamp{real: e.Real, lamport: e.Lamport, vector: e.Vector, hlc: e.HLC}
        }
        events = append(events, e)
        ready[n] = bestAt
        pos[n]++
      }
      for n := range Nodes {
        if pos[n] != len(script[n]) {
          return nil, fmt.Errorf("node %s stopped at step %d: a recv waits for a message nobody sends", NodeName(n), pos[n])
        }
      }
      return events, nil
    }
    
    // hlcLocal is the hybrid logical clock rule for a local event or a send: take the wall clock if
    // it is ahead, otherwise keep L and count up.
    func hlcLocal(h HLC, pt int) HLC {
      if pt > h.L {
        return HLC{L: pt}
      }
      return HLC{L: h.L, C: h.C + 1}
    }
    
    // hlcRecv is the rule for a receive: L becomes the largest of the node's L, the message's L and
    // the wall clock; C counts up from whichever L won.
    func hlcRecv(h, m HLC, pt int) HLC {
      l := max(h.L, m.L, pt)
      switch {
      case l == h.L && l == m.L:
        return HLC{L: l, C: max(h.C, m.C) + 1}
      case l == h.L:
        return HLC{L: l, C: h.C + 1}
      case l == m.L:
        return HLC{L: l, C: m.C + 1}
      default:
        return HLC{L: l}
      }
    }
    

    A Lamport clock jumps past every stamp it receives, so effects get larger stamps than causes. A vector clock keeps one counter per node, so it can also show that two events are concurrent.

    G

    Hybrid logical clocks

    pseudo code
    hybrid logical clockpseudo code
    on a local event or a send (pt = wall clock):
      IF pt > l:  l = pt; c = 0            // the wall clock moved on1
      ELSE:       c += 1                   // same l: count
    
    on receive of m:
      l' = max(l, m.l, pt)2
      c  = (count up from the c of whichever l won; 0 if pt won)
      l  = l'
    
    order: by l, then by c                 // l and c fit in 64 bits3
    1. 1Usually l is the wall clock itself, so the stamp reads as a time.
    2. 2A message from a node that runs ahead pulls l forward. l never goes back.
    3. 3Kulkarni et al., 2014: the stamp fits in a 64-bit NTP-format timestamp.
    Tested source Go: all four clocks in one run
    Go: all four clocks in one rungo
    // Run plays the script under a timing and a skew. Events come out in true-time order; ties go
    // to the lower node.
    func Run(tm Timing, skew int) ([]Event, error) {
      off := Offsets(skew)
      var (
        pos     [Nodes]int        // next step of each node
        ready   [Nodes]int        // true time at which each node may act next
        lamport [Nodes]int        // each node's Lamport counter
        vector  [Nodes][Nodes]int // each node's vector clock
        hlc     [Nodes]HLC        // each node's hybrid logical clock
        sent    = map[int]msgStamp{}
        events  []Event
      )
      for n := range Nodes {
        ready[n] = Start
      }
      for {
        // Pick the node whose next step can happen first. A recv cannot happen before its
        // message arrives.
        best, bestAt := -1, 0
        for n := range Nodes {
          if pos[n] >= len(script[n]) {
            continue
          }
          s := script[n][pos[n]]
          at := ready[n] + tm.Gap[n][pos[n]]
          if s.kind == KindRecv {
            m, ok := sent[s.msg]
            if !ok {
              continue // not sent yet
            }
            at = max(at, m.real+tm.Delay[s.msg])
          }
          if best < 0 || at < bestAt {
            best, bestAt = n, at
          }
        }
        if best < 0 {
          break
        }
        n, s := best, script[best][pos[best]]
        e := Event{ID: len(events), Node: n, Kind: s.kind, Value: s.value, Msg: s.msg, Peer: -1, Real: bestAt, Wall: bestAt + off[n]}
        if s.kind != KindWrite {
          e.Peer = s.peer
        }
        pt := e.Wall
    
        if s.kind == KindRecv {
          m := sent[s.msg]
          lamport[n] = max(lamport[n], m.lamport) + 1 // Lamport: jump past the sender
          for i := range Nodes {                      // vector: element-wise max
            vector[n][i] = max(vector[n][i], m.vector[i])
          }
          hlc[n] = hlcRecv(hlc[n], m.hlc, pt)
        } else {
          lamport[n]++
          hlc[n] = hlcLocal(hlc[n], pt)
        }
        vector[n][n]++
    
        e.Lamport, e.Vector, e.HLC = lamport[n], vector[n], hlc[n]
        if s.kind == KindSend {
          sent[s.msg] = msgStamp{real: e.Real, lamport: e.Lamport, vector: e.Vector, hlc: e.HLC}
        }
        events = append(events, e)
        ready[n] = bestAt
        pos[n]++
      }
      for n := range Nodes {
        if pos[n] != len(script[n]) {
          return nil, fmt.Errorf("node %s stopped at step %d: a recv waits for a message nobody sends", NodeName(n), pos[n])
        }
      }
      return events, nil
    }
    
    // hlcLocal is the hybrid logical clock rule for a local event or a send: take the wall clock if
    // it is ahead, otherwise keep L and count up.
    func hlcLocal(h HLC, pt int) HLC {
      if pt > h.L {
        return HLC{L: pt}
      }
      return HLC{L: h.L, C: h.C + 1}
    }
    
    // hlcRecv is the rule for a receive: L becomes the largest of the node's L, the message's L and
    // the wall clock; C counts up from whichever L won.
    func hlcRecv(h, m HLC, pt int) HLC {
      l := max(h.L, m.L, pt)
      switch {
      case l == h.L && l == m.L:
        return HLC{L: l, C: max(h.C, m.C) + 1}
      case l == h.L:
        return HLC{L: l, C: h.C + 1}
      case l == m.L:
        return HLC{L: l, C: m.C + 1}
      default:
        return HLC{L: l}
      }
    }
    
    at 60 ms skew
    L ran up to 58 ms ahead of a node’s own wall clock, never more than the skew.
    causal order
    Kept at every skew, over 50 seeds.

    An HLC stamp is the largest wall time seen plus a counter. It keeps causal order like Lamport and stays within the clock skew of real time.

    H

    TrueTime and commit wait

    bounded uncertainty
    true time1. TT.now() at commitearliestlatest = s2ε2. wait until TT.after(s)commit waitTT.now() moves onearliest > s3. release locks, replyEvery clock now reads later than s. A transaction that starts afterthis reply gets a larger timestamp, so timestamp order is real-time order.ε: 1 to 7 ms, 4 ms on average. Wait: at least 2ε̄, about 8 ms.
    • GPS and atomic clocks in each datacenter keep ε small.
    • The wait runs during the Paxos round, so it adds less than 2ε to a commit.
    • Without such clocks, a store bounds the skew and restarts reads that fall in the uncertainty.

    Spanner's clock returns an interval. A commit takes the top of the interval as its timestamp and waits until that time has passed everywhere, so timestamp order is real-time order.

    I

    IDs that sort by time

    layouts from each specification
    Snowflake64 bits41 ms since epoch10 node12 seqUUIDv7128 bits48 Unix ms62 random or counterULID128 bits48 Unix ms80 randomtime part: IDs sort by it first
    IDmade wheresame msclock steps back
    Snowflakenode, with a node id12-bit sequence: 4,096Wait or fail
    UUIDv7anywhererandom or counter bitsReuse last ms, or error
    ULIDanywheremonotonic: random + 1Fails on overflow
    Database sequenceone serverone counterNo clock

    Time-sorted IDs make new rows land at one end of the index. They sort by creation time only within the clock skew between generators.

    J

    Try it: the clock steps back

    recorded from the lab generator

    Read the clock. A new millisecond resets the sequence to 0, even an earlier millisecond.

    #clockID: ms · seqresult
    1100100 · 0✓ new, largest so far
    2100100 · 1✓ new, largest so far
    3101101 · 0✓ new, largest so far
    499 −3 ms99 · 0✕ below an ID already issued
    59999 · 1✕ below an ID already issued
    6100100 · 0✕ repeats an earlier ID
    7101101 · 0✕ repeats an earlier ID
    8102102 · 0✓ new, largest so far
    2 IDs repeat an earlier ID, and 2 more sort below IDs already issued. Two rows now share a primary key.

    Wait sleeps up to 5 ms; a larger step back fails at once instead.

    When the clock reads earlier than the last ID, my generator waits for a small step and fails for a large one. It never resets the sequence.

    K

    A Snowflake generator

    pseudo code
    next idpseudo code
    next():
      now = clock in ms
      IF now < last:                      // the clock stepped back1
        IF policy == fail OR last - now > max_wait2:
          RETURN error
        sleep until clock >= last
      IF now == last:
        seq += 1
        IF seq > 40953: sleep until the next ms; seq = 0
      ELSE:
        seq = 0
      last = now
      RETURN (now - epoch) << 224 | node << 12 | seq
    1. 1NTP corrected the clock. Without this check, the sequence restarts at 0 and repeats IDs.
    2. 2A large step would block for too long. Fail, alert, and take the node out.
    3. 312 bits: 4,096 IDs per ms per node. The next one waits for the next ms.
    4. 422 bits below the time: 10 for the node, 12 for the sequence.
    Tested source Go: the generator
    Go: the generatorgo
    // Next issues the next ID. Under Wait and Fail it never issues an ID at or below the last one.
    func (g *Generator) Next() (int64, error) {
      g.mu.Lock()
      defer g.mu.Unlock()
      now := g.clock.NowMS()
      if now < g.last && g.policy != Naive {
        behind := g.last - now
        if g.policy == Fail || behind > g.maxWait {
          return 0, fmt.Errorf("%w: %d ms behind the last ID", ErrClockBackwards, behind)
        }
        for now < g.last { // Wait: sleep until the clock catches up
          g.clock.Sleep(g.last - now)
          g.slept += g.last - now
          now = g.clock.NowMS()
        }
      }
      if now == g.last {
        g.seq++
        if g.seq > MaxSeq { // 4,096 IDs in this ms: wait for the next ms
          for now <= g.last {
            g.clock.Sleep(1)
            g.slept++
            now = g.clock.NowMS()
          }
          g.seq = 0
        }
      } else {
        g.seq = 0
      }
      g.last = now
      return (now-g.epoch)<<(NodeBits+SeqBits) | g.node<<SeqBits | g.seq, nil
    }
    
    • The lab issues 20,000 IDs under random backward steps of up to 5 ms: every ID is larger than the last.
    • In the trace, Wait slept 2 ms in total.

    A Snowflake ID is milliseconds, node id and sequence. The generator must never hand out a millisecond below the last one.

    L

    Capabilities used

    what each tool gives you
    toolcapabilitywhat it gives this designalso used for
    PostgresSequence: nextval()Unique, increasing numbers from one server, with no clock. Commit order can differ from number order.Surrogate keys
    PostgresWAL position of each commitOne total order of commits on one primary.Read-your-writes tokens, change feeds
    RedisINCR, one command at a timeA shared counter for IDs or versions.Rate limits, fencing tokens
    ServiceMonotonic clock in the runtimeDurations, timeouts and leases that a clock step cannot break.Latency metrics
    ServiceSnowflake, UUIDv7 or ULID made in the serviceIDs without a round trip, sorted by time within the skew.Log and event IDs
    SpannerTrueTime: now() returns an intervalCommit wait makes timestamp order match real-time order.Consistent snapshot reads
    CockroachDBHybrid logical clocks; --max-offset, 500 ms by defaultSerializable transactions without special clocks. A node that drifts too far stops itself.Reads as of a past time
    DynamoDBGlobal tables, last writer winsLoss accepted Concurrent writes in two Regions resolve by an internal timestamp.Multi-region tables
    Time syncPTP hardware clock on supported cloud instancesClock error in the microsecond range, so the skew bound in every rule above shrinks.Trading, logs

    On one database, the database orders everything. Across machines, I use logical clocks or a clock with a documented error bound.

    M

    Failure cases

    what breaks, and what saves you
    eventresultwhy it stays correctsaved by
    NTP steps a node’s clock back 3 msThe ID generator reads an earlier ms.It sleeps until the clock passes the last ID. No ID repeats.Wait policy
    The clock steps back 1 secondWaiting would block every request.The generator fails at once. An alert fires; other nodes serve.Fail policy
    A node restarts with its clock behindIts last IDs are in the future.Store the last ms used, and wait past it at start.Saved last ms
    Two generators share a node idDuplicate IDs.Node ids come from a lease, so one id has one owner.etcd lease
    Two regions write one key; clocks skewLast write wins drops the newer write.Detect the conflict with version vectors, or give the key one writer.Version vector
    A lease is checked with the wall clockA jump makes it look valid or expired.Measure the lease with the monotonic clock, and expire it early by the drift bound.Monotonic clock
    A database node drifts past the offset limitIts timestamps could break ordering.The node stops itself before it serves a wrong read.CockroachDB
    N

    Scale ladder

    start simple; climb only on a signal
    Each step adds one component1One database2+ Snowflake IDs3+ vectors or HLC4+ bounded clocksmore load →
    IDs per second, on one machine10k100k1M10MRedis INCR: 48,000 IDs per secondRedis INCR48,000Postgres nextval: 90,000 IDs per secondPostgres nextval90,000Snowflake, one node: 3,000,000 IDs per secondSnowflake, one node3,000,000Snowflake limit, one node: 4,096,000 IDs per secondSnowflake limit, one node4,096,000IDs per second, log scale
    stepaddit handlesmove up when you see
    1One database orders everything: a sequence for IDs, commit order for events.About 90,000 IDs a second from one Postgres sequence in the lab. No clock at all.The ID round trip is a bottleneck, or many services need IDs.
    2Snowflake or UUIDv7 IDs, made on each node.Up to 4,096 IDs per ms per node: 4,096,000 a second. About 3 million in the lab.Writes on several machines can conflict.
    3Version vectors, or HLC stamps on every write.Conflicts are found, not hidden. HLC stamps also serve snapshot reads by time.Transactions across regions must follow real-time order.
    4Clocks with a known error bound, plus commit wait or uncertainty restarts.Real-time order across the world, for a wait of about 2ε per commit.Top of the ladder.

    Medians of four runs, 32 clients, on an 8-core laptop shared with other jobs; runs varied by up to 2 times. The Snowflake limit is 2^12 IDs per ms × 1,000 ms.

    I let one database order everything for as long as it can. I move to per-node IDs when one counter is the limit, and to logical clocks only when writes on several machines can conflict.

    O

    Drill

    predict, then reveal

    0 of 10 known

    1. Node A’s clock is 10 ms ahead and B’s is 10 ms behind. A writes x, B reads it and writes x 5 ms later. Which value does last write wins keep?

    2. Why measure a timeout or a lease with the monotonic clock?

    3. Lamport(a) < Lamport(b). Did a happen before b?

    4. How do vector clocks find concurrent writes, and what do they cost?

    5. What does a hybrid logical clock give that a Lamport clock does not?

    6. Why does Spanner wait before it makes a commit visible?

    7. NTP steps a Snowflake node’s clock back 3 ms. What should the generator do?

    8. Two Snowflake generators get the same node id. What happens?

    9. Why are UUIDv7 or ULID better primary keys than random UUIDs?

    10. A CockroachDB node’s clock drifts too far from the others. What happens?

    P

    Numbers to say

    measured, derived or cited
    NTP
    Tens of ms over the internet; a few hundred µs on a fast LAN.
    PTP
    Microseconds, with a hardware clock on supported cloud instances.
    TrueTime
    ε from about 1 to 7 ms, 4 ms on average. Commit wait at least 2ε̄: about 8 ms.
    drift
    TrueTime applies 200 µs per second. Over a 30 s poll that is 6 ms.
    Snowflake
    41 bits of ms is 2^41 ms, about 69 years. 1,024 nodes, 4,096 IDs per ms each.
    max offset
    CockroachDB allows 500 ms by default.
    lab
    About 3 million Snowflake IDs a second; 90,000 from a Postgres sequence; 48,000 from Redis INCR.

    Cited figures from RFC 5905, the Spanner paper, the Snowflake README and the CockroachDB documentation. Lab: 8-core laptop, 32 clients.