Learn Labs
9. The Trouble with Distributed Systems

9.3 Unreliable Clocks

Lease + fencing

fencing token
15 s
lease TTL
client 1
lease t=33GC pause 20swrite
client 2
lease t=34write t=34
Problem

Split brain. Client 1's lease expired during the pause and client 2 acquired it, but client 1 does not know that and writes anyway. Two writers, corrupted data, and no error anywhere in the system.

A process-wide pause longer than the lease TTL means a client can believe it still holds the lock while another client legitimately holds it. The lock service cannot prevent this; only the storage layer can.

Eight questions applications ask of clocks — note the split:

Durations — measure an intervalPoints in time — a date and time
1. Has this request timed out yet?5. When was this article published?
2. What's the p99 response time?6. When should the reminder email be sent?
3. How many QPS in the last 5 min?7. When does this cache entry expire?
4. How long did the user spend on our site?8. What is the timestamp on this log line?
⇒ monotonic clock⇒ time-of-day clock

In a distributed system, time is a tricky business, because communication is not instantaneous. The time when a message is received is ALWAYS LATER than when it is sent — but because of variable delays, WE DON'T KNOW HOW MUCH LATER. This makes it difficult to determine THE ORDER in which things happened.

Each machine has its own clock — usually a QUARTZ CRYSTAL OSCILLATOR — not perfectly accurate, so each machine has its own notion of time. Synchronized via NTP, whose servers get time from a more accurate source such as a GPS receiver.

3.1 The two clocks

Time-of-day clockMonotonic clock
APIclock_gettime(CLOCK_REALTIME), System.currentTimeMillisclock_gettime(CLOCK_MONOTONIC/BOOTTIME), System.nanoTime
ReturnsSeconds since the epoch (midnight UTC 1 Jan 1970, Gregorian, not counting leap seconds)A meaningless absolute value — perhaps nanoseconds since boot
GuaranteeSynchronized with NTP, so a timestamp on one machine ideally means the same as on anotherGuaranteed to ALWAYS MOVE FORWARD
Can jump?YES — if the local clock is too far ahead of NTP, it may be FORCIBLY RESET and appear to JUMP BACK. Also leap seconds, and DST (avoidable by always using UTC)No jumps. NTP may SLEW it — adjust the rate at which it moves forward — by up to 0.05% by default, but cannot make it jump
Use forPoints in timeDURATIONS — timeouts, response times. "More like a stopwatch than a wall clock"
Comparable across machines?Yes (approximately)NO — the values don't mean the same thing
ResolutionHistorically coarse (10 ms steps on older Windows); less of a problem nowUsually quite good — microseconds or less

⚠️ On a server with multiple CPU sockets, there may be a SEPARATE TIMER PER CPU, not necessarily synchronized with other CPUs. The OS compensates and tries to present a monotonic view — but it is wise to take this guarantee of monotonicity WITH A PINCH OF SALT.

In a distributed system, using a monotonic clock for elapsed time is usually fine, because it doesn't assume any synchronization between nodes.

3.2 Clock synchronization is worse than you think — eight failure modes

  1. Drift. A typical quartz clock runs faster or slower than it should, Varying with temperature. Google assumes up to 200 ppm ⇒ 6 ms drift if resynced every 30 s, or A 17-second drift if resynced once a day. This Limits the best possible accuracy even if everything is working correctly.
  2. Forcible reset. If a clock differs too much from NTP, it may Refuse to synchronize or be forcibly reset. Applications observing time before and after may see Time go backward or suddenly jump forward.
  3. Silent firewalling. A node accidentally firewalled off from NTP servers may go Unnoticed for some time while drift adds up to large discrepancies. “Anecdotal evidence suggests that this does happen in practice.”
  4. Network-bounded accuracy. NTP can be only as good as the network delay. One experiment: A minimum error of 35 ms over the internet, with occasional spikes to Around a second. Large delays can cause the NTP client to Give up entirely.
  5. Wrong servers. Some NTP servers are wrong or misconfigured, Reporting time off by hours. Clients mitigate by querying several and ignoring outliers — but “it’s somewhat worrying to bet the correctness of your systems on the time you were told by a Stranger on the internet.”
  6. Leap seconds. A minute that is 59 or 61 seconds long. “The fact that leap seconds have Crashed many large systems shows how easy it is for incorrect assumptions about clocks to sneak in.” Best handling: make NTP servers “Lie” by Smearing the adjustment over a day. Leap seconds will no longer be used from 2035.
  7. VMs. The hardware clock is virtualized. When a core is shared, each VM is Paused for tens of milliseconds — which manifests to the application as The clock suddenly jumping forward. An NTP client inside the VM Doesn’t know when a pause occurs, so it may Report the clock accuracy incorrectly.
  8. Untrusted devices. On mobile/embedded devices you don’t control, you Probably cannot trust their hardware clocks at all. Some users Deliberately set an incorrect date — for example, To cheat in games.

High accuracy IS achievable if you invest: MiFID II requires high-frequency trading funds to synchronize within 100 MICROSECONDS of UTC, to help debug flash crashes and detect market manipulation. Achieved with GPS receivers and/or atomic clocks, PTP (Precision Time Protocol), and careful deployment and monitoring.

⚠️ Relying on GPS alone can be risky because GPS SIGNALS CAN EASILY BE JAMMED. In some locations (close to military facilities) this happens FREQUENTLY.

3.3 Why bad clocks are especially dangerous

Part of the problem is that INCORRECT CLOCKS EASILY GO UNNOTICED. If a CPU is defective or the network misconfigured, it most likely WON'T WORK AT ALL, so the issue is quickly spotted. If the quartz clock is defective or NTP misconfigured, MOST THINGS WILL SEEM TO WORK FINE, even as the clock drifts further from reality.

If software relies on an accurately synchronized clock, the result is more likely to be SILENT AND SUBTLE DATA LOSS than a dramatic crash.

Therefore: CAREFULLY MONITOR THE CLOCK OFFSETS between all machines. Any node whose clock drifts too far from the others should be DECLARED DEAD AND REMOVED.

3.4 Timestamps for ordering events — where it bites

node 1 (client A)node 3 (client B)node 2write x = 1 · timestamp 42.004replicated to node 3increment x → x = 2 · timestamp 42.003

The clock skew here is under 3 ms — “probably better than you can expect in practice”. Node 2 receives both writes and, under last-write-wins, keeps the one with the greater timestamp, so it keeps x = 1 and the increment is lost.

B's write is causally later than A's, but carries an earlier timestamp.

Figure 9.3.23.4 Timestamps for ordering events — where it bites

Three serious problems with client-clock LWW (Cassandra and ScyllaDB do this deliberately, to avoid the extra read round-trip needed to find the greatest existing timestamp):

  1. Database writes can mysteriously disappear. A node with a LAGGING clock is UNABLE TO OVERWRITE values previously written by a node with a FASTER clock UNTIL THE CLOCK SKEW HAS ELAPSED. This can cause ARBITRARY AMOUNTS OF DATA TO BE SILENTLY DROPPED WITHOUT ANY ERROR BEING REPORTED.
  2. LWW cannot distinguish sequential-in-quick-succession from truly concurrent writes. Additional causality tracking (VERSION VECTORS, Ch 6) is needed.
  3. Two nodes could independently generate writes with the SAME timestamp, especially at millisecond resolution. A tiebreaker (a large random number) is required — but this can ALSO lead to violations of causality.

Even with tightly NTP-synchronized clocks, you could send a packet at timestamp 100 ms (sender's clock) and have it arrive at timestamp 99 ms (recipient's clock) — SO IT APPEARS THE PACKET ARRIVED BEFORE IT WAS SENT, WHICH IS IMPOSSIBLE.

Could NTP be made accurate enough? PROBABLY NOT, because NTP's accuracy is itself limited by the NETWORK ROUND-TRIP TIME, plus quartz drift. To guarantee correct ordering you would need THE CLOCK ERROR TO BE SIGNIFICANTLY LOWER THAN THE NETWORK DELAY, WHICH IS NOT POSSIBLE.

The alternative: LOGICAL CLOCKS — based on incrementing counters rather than an oscillating quartz crystal. They do not measure the time of day or seconds elapsed, ONLY THE RELATIVE ORDERING of events. Time-of-day and monotonic clocks, which measure actual elapsed time, are PHYSICAL CLOCKS. (Ch 10.)

3.5 Clock readings with a confidence interval

It doesn't make sense to think of a clock reading as A POINT IN TIME. It is more like A RANGE OF TIMES, within a CONFIDENCE INTERVAL — a system may be 95% confident the time is between 10.3 and 10.5 seconds past the minute.

IF WE KNOW ONLY THE TIME ±100 ms, THE MICROSECOND DIGITS IN THE TIMESTAMP ARE ESSENTIALLY MEANINGLESS.

Computing the bound: with a GPS receiver or atomic clock attached, the error range is determined by the device and the signal quality. From a server: expected quartz drift since the last sync + the NTP server's uncertainty + the network round-trip time.

Unfortunately, MOST SYSTEMS DON'T EXPOSE THIS UNCERTAINTY. When you call clock_gettime, the return value doesn't tell you the expected error — SO YOU DON'T KNOW WHETHER ITS CONFIDENCE INTERVAL IS FIVE MILLISECONDS OR FIVE YEARS.

The exceptions: Google Spanner's TrueTime API and Amazon ClockBound, which return [earliest, latest].

3.6 Synchronized clocks for global snapshots — Spanner's trick

The problem: MVCC (Ch 8) requires a monotonically increasing transaction ID. On one node, a counter suffices. Distributed across many machines, a global monotonically increasing ID is difficult, because it requires COORDINATION — and the ID must reflect CAUSALITY (if B reads or overwrites a value written by A, B must have a higher ID). With lots of small, rapid transactions, creating such IDs becomes AN UNTENABLE BOTTLENECK.

Spanner's observation:

Non-overlapping intervals
AA_earliest … A_latestBB_earliest … B_latesttime →B definitely happened after A. No doubt.
Overlapping intervals
AA_earliest … A_latestBB_earliest … B_latestoverlapWe are unsure in which order A and B happened.

Spanner's solution is to deliberately wait for the length of the confidence interval before committing a read/write transaction — commit wait. This ensures that any transaction which may read the data is at a sufficiently later time that their confidence intervals do not overlap.

To keep the wait short you must keep the uncertainty small, so Google deploys a GPS receiver or an atomic clock in each datacenter, synchronizing clocks to within about 7 ms.

Figure 9.3.3Spanner's observation

The atomic clocks and GPS receivers are NOT STRICTLY NECESSARY. THE IMPORTANT THING IS TO HAVE A CONFIDENCE INTERVAL — accurate clock sources only help keep that interval SMALL.

(YugabyteDB can leverage ClockBound on AWS, and several other systems now rely on clock synchronization to various degrees.)


On this page