A compare-and-set needs an owner, so Cassandra elects one for each write

In "Lightweight transactions in Cassandra 2.0" [1], Jonathan Ellis introduced Cassandra's lightweight transactions with account registration: two people try to claim the same login at the same moment, and exactly one should get it. Each client checks that the login is free, then inserts it. Done as two steps, that's the classic check-then-act race; registration needs a compare-and-set, which checks and writes as one atomic step.

Cassandra already had a strong-sounding tool. With three replicas, a QUORUM write waits for two to acknowledge, and a QUORUM read asks two. Since 2 + 2 > 3, every read overlaps every acknowledged write, which is why this is often called strong consistency. So why isn't a QUORUM read followed by a QUORUM write a safe compare-and-set?

Because the overlap guarantees visibility, not atomicity. It promises that a read sees every write acknowledged before the read started. Registration is two operations, and two clients can interleave them.

A check-then-act race, and compare-and-set with an ownerTwo timelines of clients A and B registering the same login. With a quorum read then a quorum write, A checks and sees the login free, B checks and also sees it free, then both write, and both think they won. With one owner, A's check and write run as a single unit; B waits, then checks, sees the login taken, and is rejected.Quorum read, then writeboth checks run before either writeAfreewritethinks it wonBfreewritethinks it wonOne ownereach compare-and-set runs whole, one at a timeAfreewritegets the loginBwaitstakenrejected
  • check (what it saw)
  • write
Quorums make the latest write visible, but two clients can both check before either writes. An owner runs each compare-and-set as one step.

A compare-and-set needs an owner: one place that puts the operations on a key in order. In Postgres the owner is the primary. Replicas don't take writes, and with a unique index on the login the second insert waits for the first to commit, then fails. In Cassandra any replica accepts a write for any key, so no machine owns the row.

So IF NOT EXISTS makes an owner for one operation. The coordinator runs a Paxos round on the partition: it asks the replicas to promise to ignore any lower-numbered ballot, and once a majority has promised, it is in effect the partition's owner. Then it reads the current value and writes only if the check passes.

Postgres pays for its owner too, just once. Choosing a primary is itself an agreement: a failover has to make sure only one primary survives. After that, the owner orders every write cheaply. Cassandra is leaderless by design, so it elects an owner for each conditional write instead.

References

  1. Lightweight transactions in Cassandra 2.0 [link]
    Ellis, J., 2013. DataStax Blog.