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.

Advertisement

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.

  1. Local event on process i: V[i] += 1.
  2. Send on process i: V[i] += 1, then attach a copy of V to the message.
  3. 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.

Advertisement

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.

Three processes, message passing, and the vector each event carriesP1P2P3a [1,0,0]b send [2,0,0]c [0,1,0]d recv [2,2,0]f send [2,3,0]e [0,0,1]g recv [2,3,2]m1m2a precedes g: every entry of [1,0,0] is at most the matching entry of [2,3,2], and one is smallerb and e are concurrent: [2,0,0] has the larger P1 entry, [0,0,1] the larger P3 entryReceive rule: take the element-wise max with the message's vector, then add one to your own entry
Each event is stamped with its process's vector after applying the rules. Comparing two stamps entry by entry decides whether one event could have influenced the other.

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.

MechanismEntry perIncremented onUsed for
Lamport clockOne scalarEvery event; max+1 on receiveTotal order consistent with causality; cannot detect concurrency
Vector clockProcessEvery event, including sends and receivesDeciding happens-before between any two events, debugging, causal delivery
Version vectorReplicaOnly when a replica accepts an updateComparing versions of one data item; detecting conflicting writes
Dotted version vectorReplica, plus a dot on each siblingUpdate; the dot names the single event that created a valueServer-side version tracking without sibling explosion
Hybrid logical clockOne timestampPhysical time plus a logical counterSnapshot 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], context

This 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

FailureWhat you seeCauseFix
False conflictSiblings for writes that were really sequentialClient wrote without sending the context it read; or an entry was prunedAlways round-trip the context; prune conservatively
Silent lost updateAn acknowledged write disappearsServer compared with the wrong context, or fell back to timestampsTest with concurrent writers; never let a missing context mean overwrite
Sibling explosionHundreds of values under one key, slow readsMany clients writing through one replica with plain version vectors, or nobody resolving siblingsDotted version vectors; resolve on read; alert on sibling count
Resurrected deletesRemoved items come back after a mergeUnion merge with no tombstonesTombstones 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

  1. Implement the VectorClock class above and reproduce the three-process example; check every pairwise comparison by hand.
  2. Decide whether your use case needs vector clocks (events) or version vectors (data versions), and write that down in the design.
  3. Make the causal context a required parameter of every write in your client library, including retries.
  4. 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.
  5. Choose a merge function per data type; use tombstones or an observed-remove set wherever items can be deleted.
  6. Add dashboards for sibling count per read, vector size and context-free writes, and define a pruning policy that only removes old, converged entries.
  7. 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.
Key takeaway: A vector clock records, per participant, how many of its events you know about, and comparing two vectors entry by entry tells you exactly whether one event could have influenced another or whether they were concurrent. Use version vectors for data, round-trip the causal context on every write, keep concurrent writes as siblings and merge them with type-aware logic, and watch vector size and sibling counts. When you cannot merge, prevent concurrency with consensus; when you only need an order, a hybrid logical clock is cheaper.