Two replicas of a shopping cart each accept a write during a network partition. When the partition heals, the system has two values for one key and must answer a simple question: did one write see the other, or were they made independently? If one saw the other, the later value wins and nothing is lost. If they were independent, throwing either away loses a customer's action. Wall-clock timestamps cannot answer the question reliably. Vector clocks can.
This article builds vector clocks from first principles: happens-before, the exact rules, a worked example, version vectors and dotted version vectors, siblings in a replicated store, and the operational problems of metadata growth and pruning.
First principles: why time cannot order events
Every machine has its own clock, and clocks disagree. NTP typically keeps servers within milliseconds of each other, but a virtual machine pause, a leap-second smear or a misconfigured host can put one node seconds away from the rest. A timestamp order therefore says nothing reliable about which write came first, let alone whether one writer knew about the other.
Leslie Lamport's 1978 paper replaced time with causality. Event x happens before event y, written x → y, if x came first on the same process, or x sends a message that y receives, or a chain of such steps links them. If neither x → y nor y → x, the events are concurrent: neither could have influenced the other. Happens-before is exactly the information a replicated system needs, because it separates overwrites (the writer saw the old value) from conflicts (the writer did not).
Lamport clocks, a single counter per process, give a timestamp L such that x → y implies L(x) < L(y). The converse fails: L(x) < L(y) tells you nothing about whether x influenced y. Vector clocks, introduced independently by Fidge and Mattern in 1988, keep one counter per process and restore the missing direction: x → y if and only if V(x) < V(y). Charron-Bost later showed that, in general, no characterisation of causality with fewer than n entries exists for n processes, so the size is not an implementation accident.
The rules
Each process i keeps a vector V of n non-negative integers, all zero at start. V[j] means: the number of events of process j that process i knows about. Three rules maintain it.
- Local event on process i: V[i] += 1.
- Send on process i: V[i] += 1, then attach a copy of V to the message.
- Receive of a message carrying W on process i: for every j, V[j] = max(V[j], W[j]); then V[i] += 1.
Comparison is element-wise. V ≤ W when every entry of V is at most the matching entry of W. V < W when V ≤ W and V ≠ W; then the event stamped V happened before the event stamped W. If neither V ≤ W nor W ≤ V, the events are concurrent. Here is a complete implementation using dictionaries, so absent entries mean zero and the set of processes can grow.
from enum import Enum
class Order(Enum):
BEFORE = "before"
AFTER = "after"
EQUAL = "equal"
CONCURRENT = "concurrent"
class VectorClock:
def __init__(self, entries=None):
self.v = dict(entries or {})
def tick(self, node): # local event or send
self.v[node] = self.v.get(node, 0) + 1
return self
def merge(self, other): # element-wise max
for node, count in other.v.items():
if count > self.v.get(node, 0):
self.v[node] = count
return self
def receive(self, node, other): # merge, then tick
return self.merge(other).tick(node)
def compare(self, other):
keys = self.v.keys() | other.v.keys()
le = all(self.v.get(k, 0) <= other.v.get(k, 0) for k in keys)
ge = all(self.v.get(k, 0) >= other.v.get(k, 0) for k in keys)
if le and ge:
return Order.EQUAL
if le:
return Order.BEFORE
if ge:
return Order.AFTER
return Order.CONCURRENT
def copy(self):
return VectorClock(self.v)Comparison and message overhead are both O(n) in the number of entries: exact causality costs metadata proportional to the participants.
Worked example: three processes
Follow the diagram. P1 has a local event a, giving [1,0,0], then sends m1 at event b with [2,0,0]. Meanwhile P2 has a local event c, [0,1,0]. When m1 arrives at P2, the receive rule takes the element-wise max of [0,1,0] and [2,0,0], which is [2,1,0], then increments P2's own entry, giving d = [2,2,0]. P2 sends m2 at event f, [2,3,0]. P3 had a local event e, [0,0,1]; on receiving m2 it takes max([0,0,1],[2,3,0]) = [2,3,1] and increments its entry: g = [2,3,2].
Now ask questions. Did a happen before g? [1,0,0] is at most [2,3,2] everywhere and smaller somewhere, so yes: the chain a, b, m1, d, f, m2, g carried the knowledge. Is b related to e? [2,0,0] against [0,0,1]: each has an entry larger than the other, so they are concurrent. Is c before g? [0,1,0] < [2,3,2], yes, through P2's own later send. Is c related to a? Concurrent.
Vector clocks, version vectors and dotted version vectors
The names are often used interchangeably, which causes bugs. A vector clock orders events: every send, receive and local step increments it. A version vector orders versions of a data item: a replica increments its entry only when it accepts an update to that item, and merging two replicas' states takes the element-wise max without incrementing, because synchronising is not a new update. Databases want version vectors; debuggers, causal message delivery and distributed snapshot algorithms want vector clocks.
Concurrent does not mean conflicting. Concurrent increments of a counter, or writes to different fields, merge mechanically; whether concurrency is a problem depends on the data type, which is why CRDTs pair naturally with version vectors.
| Mechanism | Entry per | Incremented on | Used for |
|---|---|---|---|
| Lamport clock | One scalar | Every event; max+1 on receive | Total order consistent with causality; cannot detect concurrency |
| Vector clock | Process | Every event, including sends and receives | Deciding happens-before between any two events, debugging, causal delivery |
| Version vector | Replica | Only when a replica accepts an update | Comparing versions of one data item; detecting conflicting writes |
| Dotted version vector | Replica, plus a dot on each sibling | Update; the dot names the single event that created a value | Server-side version tracking without sibling explosion |
| Hybrid logical clock | One timestamp | Physical time plus a logical counter | Snapshot reads and ordering near wall-clock time; not concurrency detection |
Architecture: siblings in a replicated key-value store
Dynamo-style stores put version vectors at the centre of the write path. It reads a key, receives the value or values plus an opaque causal context (the version vector), and sends that context back with its next write. The coordinating replica increments its own entry on top of the context and stores the result. On each replica, a new version replaces any stored version whose vector it dominates and is kept alongside any version it is concurrent with. Those concurrent values are siblings.
Walk through a cart. A client writes {milk} through replica A; the version is {A:1}. The same client adds eggs through A with context {A:1}, producing {A:2}, which replaces it. During a partition a second client that had read {A:1} adds bread through B; B stores {A:1, B:1} = {milk, bread}. When anti-entropy or read repair brings the replicas together, {A:2} and {A:1, B:1} are concurrent, so both are kept. The next reader receives two siblings and the merged context {A:2, B:1}. It resolves them, here with a set union to {milk, eggs, bread}, and writes through replica C with that context, producing {A:2, B:1, C:1}, which dominates both siblings and replaces them everywhere.
Two details decide whether this works. First, resolution is the application's job, because only it knows that a cart merges by union while a bank balance does not. Second, a union merge resurrects deletes: if one sibling removed an item and the other did not, the union brings it back. Removal needs tombstones or an observed-remove set.
def put(store, key, value, context, coordinator):
"""Write through a coordinating replica using version vectors."""
new_vv = context.copy().tick(coordinator)
kept = []
for sibling_value, sibling_vv in store.get(key, []):
order = sibling_vv.compare(new_vv)
if order in (Order.BEFORE, Order.EQUAL):
continue # superseded by the new write
kept.append((sibling_value, sibling_vv)) # concurrent: keep
kept.append((value, new_vv))
store[key] = kept
return new_vv
def get(store, key):
siblings = store.get(key, [])
context = VectorClock()
for _, vv in siblings:
context.merge(vv) # the context covers every sibling
return [v for v, _ in siblings], contextThis sketch has a known flaw. If two clients both write through replica A with the same stale context {A:1}, the sketch stamps both {A:2}, the second compares EQUAL to the first and replaces it. A replica that tracks its own counter would issue {A:2} and then {A:3} instead, and {A:3} dominates {A:2}. Either way the first write is silently discarded even though its author never saw the second. Plain per-replica version vectors cannot describe two concurrent values created by the same replica. Riak adopted dotted version vectors, from work by Preguiça, Baquero and colleagues: each stored value carries a dot, the single (replica, counter) event that created it, plus the context it was written against, so the server can tell that the value at dot (A,3) was written with context {A:1} and does not cover the value at dot (A,2). If you implement server-side versioning, use dots or per-client entries.
Keeping the metadata bounded
The vector has one entry per participant, so choosing the participant is a design decision. Entries per client are correct but grow with every client that ever wrote the key. Entries per replica keep the vector to roughly the replication factor for keys that stay on their home replicas, but sloppy quorums and hinted handoff add entries for fallback nodes, and cluster churn leaves entries for nodes that no longer exist.
The Dynamo paper describes truncation: when the number of (node, counter) pairs reaches a threshold (the paper gives 10 as an example), the oldest pair by an attached wall-clock timestamp is removed. A pruned entry forgets that some writes happened, so a later comparison can report false concurrency or, worse, make an old version look dominant. Pruning is only safe when every replica has already converged past the pruned entries.
Key entries on stable replica identities rather than pods that come and go, prune only old entries, and log every pruning event.
Failure modes
| Failure | What you see | Cause | Fix |
|---|---|---|---|
| False conflict | Siblings for writes that were really sequential | Client wrote without sending the context it read; or an entry was pruned | Always round-trip the context; prune conservatively |
| Silent lost update | An acknowledged write disappears | Server compared with the wrong context, or fell back to timestamps | Test with concurrent writers; never let a missing context mean overwrite |
| Sibling explosion | Hundreds of values under one key, slow reads | Many clients writing through one replica with plain version vectors, or nobody resolving siblings | Dotted version vectors; resolve on read; alert on sibling count |
| Resurrected deletes | Removed items come back after a merge | Union merge with no tombstones | Tombstones or an observed-remove set CRDT |
The most common production bug is a client that writes without the context it read, often from a retry path that rebuilt the request. Such a write is concurrent with every stored version and becomes one more sibling, or, where a missing context means blind overwrite, silently wins. Make the context a required argument of every write API.
Trade-offs and alternatives
Vector clocks buy exact concurrency detection and pay for it in metadata, client complexity and application-level merge logic. Cassandra resolves conflicting cells by last-writer-wins on timestamps, accepting that a concurrent write can be lost, in exchange for small fixed metadata and no siblings. Systems that need snapshot reads close to real time use hybrid logical clocks, which fit in 64 bits and preserve happens-before in one direction, like Lamport clocks, but cannot tell you that two events were concurrent.
Use version vectors when losing a concurrent write is unacceptable and the data can be merged by the application or a CRDT: carts, collaborative documents, device sync, offline-first mobile data. Use consensus instead when the data cannot be merged at all, such as uniqueness constraints and balances, because then you need to prevent concurrency rather than detect it. And if you need causal ordering for reads across many keys rather than versions of one key, read causal consistency, which builds on the same ideas at the level of a whole datastore.
What to do next
- Implement the VectorClock class above and reproduce the three-process example; check every pairwise comparison by hand.
- Decide whether your use case needs vector clocks (events) or version vectors (data versions), and write that down in the design.
- Make the causal context a required parameter of every write in your client library, including retries.
- Use dotted version vectors or per-client entries if many clients write through the same replica, and add a concurrent-writer test that would catch the {A:2}/{A:3} bug.
- Choose a merge function per data type; use tombstones or an observed-remove set wherever items can be deleted.
- Add dashboards for sibling count per read, vector size and context-free writes, and define a pruning policy that only removes old, converged entries.
- Write a property test that replays random interleavings of writes from several clients through several replicas and asserts that no write is discarded unless its author's context was covered.