Learn Labs
6. Replication

6.5 Detecting Concurrent Writes

Conflicts might be detected as the writes happen — but not always; they could also be detected later, during read repair, hinted handoff, or anti-entropy.

client Aset X = Aclient Bset X = Bnode 1X = Anode 2X = Bnode 3X = ANode 1 never receives B (transient outage). Node 2 gets A then B; node 3 gets B then A.If each node simply overwrote on each write, the replicas stay permanently inconsistent.
Figure 6.5.1Detecting Concurrent Writes

To become eventually consistent, replicas must converge — using any of §3.4's mechanisms: LWW (Cassandra, ScyllaDB), manual resolution, or CRDTs (Riak).

LWW is easy to implement. But a timestamp DOESN'T TELL YOU WHETHER TWO VALUES ARE ACTUALLY CONFLICTING (written concurrently) or not (written one after another). To resolve conflicts explicitly, the system must take more care to detect concurrent writes.

5.1 The happens-before relation

An operation A HAPPENS BEFORE another operation B if B knows about A, or depends on A, or builds upon A in some way.

Two operations are CONCURRENT if neither happens before the other.

Three possibilities for any A and B: A happened before B, B happened before A, or they are concurrent.

  • §3.2's overtaking example: the insert happens before the increment, because the value incremented by B is the value inserted by A — B builds upon A, so B must have happened later. B is causally dependent on A.
  • §5's example: the writes are concurrent — when each client starts, it doesn't know another client is also operating on that key. There is no causal dependency.

If one operation happened before another, the later one should OVERWRITE the earlier. If the operations are concurrent, we have a CONFLICT that needs to be resolved.

Concurrency, time, and relativity:

It is NOT important whether they literally overlap in time. Because of clock problems, it's quite difficult to tell whether two things happened at exactly the same time (Ch 9). For defining concurrency, exact time doesn't matter — two operations are concurrent if they are both UNAWARE OF EACH OTHER, regardless of physical time.

People connect this to special relativity: information cannot travel faster than light, so two events some distance apart cannot affect each other if the time between them is shorter than light's travel time. In computer systems, two operations might be concurrent EVEN THOUGH the speed of light would in principle have allowed one to affect the other — if the network was slow or interrupted, two operations can occur some time apart and still be concurrent, because network problems prevented one from knowing about the other.

5.2 The algorithm (single replica)

  1. The server maintains a version number per key, increments it on every write, and stores the new version number along with the value.
  2. On read, the server returns ALL SIBLINGS — all values not overwritten — plus the latest version number. A client must READ a key BEFORE writing.
  3. On write, the client must include the version number from the prior read, and must MERGE together all values it received in that read. The write response also returns all siblings, allowing several writes to be chained.
  4. On receiving a write with a particular version number, the server can OVERWRITE all values with that version number OR BELOW (it knows they've been merged into the new value) but must KEEP all values with a HIGHER version number (those are concurrent with the incoming write).

The server can determine whether two operations are concurrent JUST BY LOOKING AT VERSION NUMBERS. It does not need to interpret the value itself — so the value could be any data structure.

A write WITHOUT a version number is concurrent with all other writes, so it will not overwrite anything — it will just be returned as one of the values on subsequent reads.

5.3 The shopping cart trace — follow this carefully

StepClientActionSent (value, version)Server state after
11add milk[milk], —v1: [milk]
22add eggs, unaware of milk[eggs], —v1: [milk]
v2: [eggs] — siblings; the server returns both plus v2
31add flour, unaware of eggs. Client had v1 = [milk], so v1 is overwritten and v2 [eggs] is kept (2 > 1, concurrent)[milk, flour], v1v2: [eggs]
v3: [milk, flour]
42add ham. Client had received [milk] and [eggs] at v2, merges them and adds ham, so v2 is overwritten and v3 kept (3 > 2, concurrent)[eggs, milk, ham], v2v3: [milk, flour]
v4: [eggs, milk, ham]
51add bacon. Client had [milk, flour] and [eggs] at v3, merges and adds bacon, so v3 is overwritten and v4 kept (4 > 3, concurrent)[milk, flour, eggs, bacon], v3v4: [eggs, milk, ham]
v5: [milk, flour, eggs, bacon]

In this example, the clients are NEVER fully up to date with the data on the server, since there is always another operation going on concurrently. But old versions DO get overwritten eventually, and NO WRITES ARE LOST.

5.4 Version vectors

A single version number is not sufficient when there are MULTIPLE REPLICAS accepting writes concurrently.

Instead, use a version number PER REPLICA as well as per key. Each replica increments its own version number when processing a write, and also keeps track of the version numbers it has seen from each of the other replicas. This information indicates which values to overwrite and which to keep as siblings.

The collection of version numbers from all the replicas is called a VERSION VECTOR.

The most interesting variant is the DOTTED VERSION VECTOR, used in Riak 2.0.

Like the single version numbers, version vectors are sent from replicas to clients on read and must be sent back on write. (Riak encodes it as a string it calls causal context.) The version vector lets the database distinguish between overwrites and concurrent writes.

The version vector also ensures that it is SAFE TO READ FROM ONE REPLICA AND SUBSEQUENTLY WRITE BACK TO ANOTHER. Doing so may create siblings, but NO DATA IS LOST as long as siblings are merged correctly.

⚠️ Version vector ≠ vector clock. They're sometimes conflated. The difference is subtle — in brief, when comparing the STATE OF REPLICAS, version vectors are the right data structure to use.


On this page