TT Lab
Get started
Learn Learning paths Courses

CS for Building Good Services — Relearning Textbook Ideas by Measuring

Numbers Keep Causality; Vectors Recognize Concurrency

Continue in TT Lab

In one line

What decides "what came first" in a distributed system is not the wall clock but the messages. A Lamport clock attaches numbers that do not go against causality, and a vector clock goes one step further and also tells you that "the two events do not know about each other (they are concurrent)." If you do not know this difference, LWW quietly drops updates, and the retries attached at each layer multiply.

Why this was needed

It is common to decide an order by looking at the timestamps of two log lines. But if the two lines came from different servers, that comparison rests on the assumption that "the two clocks are right." Why and by how much clocks drift apart, and NTP, are covered as an incident by the "The Logs Came From the Future — Five Incidents a Clock Made" course. This module starts from the opposite question: what can you say for certain about order without synchronizing the clocks?

In Time, Clocks, and the Ordering of Events in a Distributed System (CACM vol. 21, no. 7, 1978), Lamport defined "happened before" (→) not through physical time but only through events observable inside the system. There are three conditions. If a comes before b in the same process, then a → b. If a is the sending of a message and b is the receipt of that message, then a → b. If a → b and b → c, then a → c. And two distinct events for which neither a → b nor b → a holds are called concurrent. The paper stresses that this relation is only a partial order on the set of events. It means that pairs whose order is not determined exist from the start, and the paper notes that problems arise because people are not sufficiently aware of this fact.

How it works

The Lamport clock. There are two implementation rules. IR1: each process increments its clock between two successive events. IR2: a message carries the sending time Tm, and the receiver sets its own clock to a value "greater than or equal to its current value and greater than Tm." A common implementation is max(자기 값, Tm) + 1 on receipt (the larger of the clock's own value and Tm, plus 1). If you drop the max and only add +1, a receipt can get a number smaller than its send, and the rule breaks.

What this rule guarantees is one thing, the Clock Condition: if a → b then C(a) < C(b). The paper states clearly that the converse cannot be expected. For the converse to hold, two concurrent events would both have to be at the same time, which is impossible. So if you see C(a) < C(b) and say a → b, you are wrong. The paper breaks ties between equal numbers with an arbitrary order assigned to the processes to make a total order, but that is only one of many line-ups that do not contradict causality, and it does not tell you about causality.

The vector clock. If you need the converse too, one number is not enough. Virtual Time and Global States of Distributed Systems (1989) by Mattern uses a vector with as many slots as there are processes. The same paper notes that Fidge independently came up with the same idea. At each event you increment your own slot by 1, send the vector along with messages, and the receiver merges by taking the max slot by slot. Comparison is also done slot by slot. If every slot is less than or equal and the two differ, then u < v, and if neither holds, they are concurrent (u ‖ v). Theorem 10 of the paper says that e < e′ and C(e) < C(e′) are equivalent, and that e ‖ e′ and C(e) ‖ C(e′) are also equivalent.

Lamport clock Vector clock
If a → b L(a) < L(b) V(a) < V(b)
If the value is smaller You know nothing a → b
Can it recognize concurrency No V(a) ‖ V(b)
Size One integer As many as the number of processes

The price is size. With n processes, every message carries n slots, and in a system where participants come and go, managing the slots becomes a separate job.

What it looks like in the field

LWW drops updates. The Cassandra documentation says that unlike the original Dynamo paper, which reconciled concurrent updates with vector clocks, Cassandra uses the simpler last-write-wins. Every change is stamped with a timestamp from the client's or the coordinator's clock, and the latest value wins. The same document says correctness depends on that clock, so you must run synchronization such as NTP. There are two layers of loss hidden here. One of two concurrent writes disappears with no error, and the write of a replica with a fast clock beats the write of a slow replica that was made after seeing it (a causal inversion). If you attach vector clocks, at least the fact that "these two were concurrent" remains, and the application can decide whether to merge them or ask a person.

Retries multiply. The chapter Addressing Cascading Failures of the Google SRE book warns that if you retry at several layers, one top-level request can generate as many calls to the lowest layer as the product of the per-layer attempt counts. In the book's example, when the DB is overloaded and the backend, frontend, and JavaScript each retry 3 times (4 attempts), a single user action reaches the DB 64 times, and it does so precisely when the DB is weakest. At which layer to retry, plus retry budgets, jitter, and idempotency keys, is covered as design by the "Idempotency — Two Clicks, One Charge" course, and terms, another device that orders things by number, are covered by the "There Were Two Leaders" course.

What you will do in the next lab

You will attach Lamport clocks to the event records of four processes, judge happened-before by its definition, and see for which pairs "a smaller number means earlier" is wrong. You will count all concurrent pairs with vector clocks, and find the list of conflicts and the updates that LWW silently drops in writes to the same key on three replicas whose wall clocks are skewed. Finally you will calculate how many times the DB calls multiply under a per-layer retry policy and write a summary report.