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.
- 1Wall clocks disagree by milliseconds or more. Never order causally related events by wall time.
- 2Measure durations with the monotonic clock; use the wall clock only to tell the time.
- 3Lamport clocks keep causal order; vector clocks also detect concurrency; HLC stays near wall time.
- 4Time-sorted IDs must wait or fail when the clock steps back, or they repeat.
- 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.
| clock | keeps causal order | finds concurrency | near real time |
|---|---|---|---|
| Wall clock | No | No | Within the skew |
| Monotonic clock | One machine | No | Arbitrary start |
| Lamport | Yes | No | No |
| Vector | Yes | Yes | No |
| Hybrid logical (HLC) | Yes | No | Within the skew |
| TrueTime interval | With commit wait | No | Within ε |
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.
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.
- B: recv m1 97
- B: b1 99
- B: send m2 103
- C: c1 104
- C: send m3 110
- A: a1 114
- A: send m1 115
- C: recv m2 119
- A: recv m3 123
- C: c2 123
- C: send m4 127
- A: recv m4 142
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.
- 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.
| use | clock | status |
|---|---|---|
| Show the time, stamp a log line | wall | Approved |
| Timeout, lease, latency | monotonic | Approved |
| Timeout from two wall readings | wall | Not approved |
| Order events across machines | wall | Not approved |
| Compare monotonic readings of two machines | monotonic | Different starts |
A timeout, a lease or a latency is a duration, so I measure it with the monotonic clock. The wall clock can jump.
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- 1The receive gets a larger stamp than the send. Lamport, 1978: if a happens before b, C(a) < C(b).
- 2The node now knows everything the sender knew.
- 3The lab checks this rule against the true causal order for every pair, over 50 seeds and 4 skews.
- 4Concurrent: a conflict for last write wins to hide, or for the app to merge.
Tested source Go: all four clocks in one run
// 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.
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- 1Usually l is the wall clock itself, so the stamp reads as a time.
- 2A message from a node that runs ahead pulls l forward. l never goes back.
- 3Kulkarni et al., 2014: the stamp fits in a 64-bit NTP-format timestamp.
Tested source Go: all four clocks in one run
// 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.
- 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.
| ID | made where | same ms | clock steps back |
|---|---|---|---|
| Snowflake | node, with a node id | 12-bit sequence: 4,096 | Wait or fail |
| UUIDv7 | anywhere | random or counter bits | Reuse last ms, or error |
| ULID | anywhere | monotonic: random + 1 | Fails on overflow |
| Database sequence | one server | one counter | No 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.
Read the clock. A new millisecond resets the sequence to 0, even an earlier millisecond.
| # | clock | ID: ms · seq | result |
|---|---|---|---|
| 1 | 100 | 100 · 0 | ✓ new, largest so far |
| 2 | 100 | 100 · 1 | ✓ new, largest so far |
| 3 | 101 | 101 · 0 | ✓ new, largest so far |
| 4 | 99 −3 ms | 99 · 0 | ✕ below an ID already issued |
| 5 | 99 | 99 · 1 | ✕ below an ID already issued |
| 6 | 100 | 100 · 0 | ✕ repeats an earlier ID |
| 7 | 101 | 101 · 0 | ✕ repeats an earlier ID |
| 8 | 102 | 102 · 0 | ✓ new, largest so far |
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.
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- 1NTP corrected the clock. Without this check, the sequence restarts at 0 and repeats IDs.
- 2A large step would block for too long. Fail, alert, and take the node out.
- 312 bits: 4,096 IDs per ms per node. The next one waits for the next ms.
- 422 bits below the time: 10 for the node, 12 for the sequence.
Tested source Go: the generator
// 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.
| tool | capability | what it gives this design | also used for |
|---|---|---|---|
| Postgres | Sequence: nextval() | Unique, increasing numbers from one server, with no clock. Commit order can differ from number order. | Surrogate keys |
| Postgres | WAL position of each commit | One total order of commits on one primary. | Read-your-writes tokens, change feeds |
| Redis | INCR, one command at a time | A shared counter for IDs or versions. | Rate limits, fencing tokens |
| Service | Monotonic clock in the runtime | Durations, timeouts and leases that a clock step cannot break. | Latency metrics |
| Service | Snowflake, UUIDv7 or ULID made in the service | IDs without a round trip, sorted by time within the skew. | Log and event IDs |
| Spanner | TrueTime: now() returns an interval | Commit wait makes timestamp order match real-time order. | Consistent snapshot reads |
| CockroachDB | Hybrid logical clocks; --max-offset, 500 ms by default | Serializable transactions without special clocks. A node that drifts too far stops itself. | Reads as of a past time |
| DynamoDB | Global tables, last writer wins | Loss accepted Concurrent writes in two Regions resolve by an internal timestamp. | Multi-region tables |
| Time sync | PTP hardware clock on supported cloud instances | Clock 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.
| event | result | why it stays correct | saved by |
|---|---|---|---|
| NTP steps a node’s clock back 3 ms | The 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 second | Waiting would block every request. | The generator fails at once. An alert fires; other nodes serve. | Fail policy |
| A node restarts with its clock behind | Its 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 id | Duplicate IDs. | Node ids come from a lease, so one id has one owner. | etcd lease |
| Two regions write one key; clocks skew | Last 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 clock | A 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 limit | Its timestamps could break ordering. | The node stops itself before it serves a wrong read. | CockroachDB |
| step | add | it handles | move up when you see |
|---|---|---|---|
| 1 | One 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. |
| 2 | Snowflake 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. |
| 3 | Version 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. |
| 4 | Clocks 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.
0 of 10 known
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?
Why measure a timeout or a lease with the monotonic clock?
Lamport(a) < Lamport(b). Did a happen before b?
How do vector clocks find concurrent writes, and what do they cost?
What does a hybrid logical clock give that a Lamport clock does not?
Why does Spanner wait before it makes a commit visible?
NTP steps a Snowflake node’s clock back 3 ms. What should the generator do?
Two Snowflake generators get the same node id. What happens?
Why are UUIDv7 or ULID better primary keys than random UUIDs?
A CockroachDB node’s clock drifts too far from the others. What happens?
- 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.