Distributed systems
Part 6 of 6 · CRDTsProduction CRDTs & When NOT to Use Them (Riak, Redis CRDT, collab apps, vs Raft/linearizability)
Riak data types, Redis Active-Active, and Yjs or Automerge are production CRDTs. Unique names, non-negative money, and exactly-once effects still belong on a linearizable store.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
How do you count likes in several regions?
Answer
A G-Counter or PN-Counter, either in Redis Active-Active or in the application. Roll analytics up downstream. Do not use the same counter as a wallet balance.
L2
How do you build collaborative notes with a paid export?
Answer
Yjs or Automerge for the draft, snapshots for durability, and a publish path that writes one immutable version and charges through a ledger. The CRDT does not send the invoice.
L3
What is wrong with Active-Active inventory when stock is 1?
Answer
Two regions can both decrement. The merge converges on zero or negative and both customers were told yes. Reserve in one linearizable store, or accept oversell and a compensation.
L4
Is a Cosmos DB last-writer-wins policy a CRDT?
Answer
Last-writer-wins is a register CRDT. A custom conflict procedure might be anything. Neither statement means the account has an observed-remove set. Name the policy.
L5
How do CRDTs interact with an MVCC database?
Answer
The database versions the row that holds the blob, which gives you local durability and a snapshot of that row. Concurrent writers on other replicas still need a CRDT merge. Snapshot isolation will not combine two JSON edits.
L6
Why did Active-Active memory grow?
Answer
Dots, tombstones, version vectors, and uncompacted deltas. Churn plus a replica that has not acked. The value the user sees can be small while the metadata is not.
L7
When is two-phase commit the wrong rescue?
Answer
When the participants are HTTP APIs or the path is cross-region. The CRDT was the wrong tool for the invariant, and two-phase commit is often the wrong strong tool. Prefer one database you can commit, or a saga if the business already accepts compensation.
Failure modes
Geo decrement of the last unit
Both regions apply a local decrement. The CRDT converges. Two orders were accepted. The datatype did not know stock was an invariant.
Two successful creates of one username
Each region added the name to its set or wrote its row. Strong eventual consistency keeps both until a later janitor notices. The primary key needed a reject.
Publish fired from every replica
The draft merged, and each replica treated the merge as the moment to charge or to email. Exactly-once was an assumption about the stream.
Tombstones with no budget
Active-Active or a document CRDT retains remove metadata until stability. Memory and sync bytes climb. The feature looks idle.
Misconceptions
We use Redis, so we are eventually consistent.
A single primary is linearizable on that primary. Cache-aside is a fill in front of a database. Active-Active is the CRDT mode. Name one.
If CRDTs are not strong enough, wrap the path in two-phase commit.
Most of these invariants want one database transaction in one region. Two-phase commit blocks prepared participants and is a poor wide-area default.
A CRDT document enforces who may edit it.
Merge is not authorization. A peer who can submit updates can submit hostile ones. Sign updates, terminate them on a server you trust, or both.
Interviewer traps
Walk through Raft log matching to justify the like button.
The like is a counter. Say so, then say the username create is the Raft or unique-index path.
Describe a saga for every CRDT field.
A saga compensates a business step that already committed locally. It is the checkout across payment, stock, and shipping. It is not the merge function for a paragraph.
Design scenario
Same prompt for every reader.
Requirements
Likes and the note converge. The cart converges and is rechecked. Stock and payment each happen once. A banned account stays banned.
Traffic / scale
Likes and note edits dominate. Checkout is hundreds per second, routed to one inventory region.
Latency
Note keystrokes stay local. Checkout waits on one database round trip plus the payment call.
Consistency
Strong eventual consistency on the note, the likes, and the cart. Linearizability on stock reservation and on the payment id.
Availability
A partitioned region still records likes and cart edits. It stops confirming checkout of scarce stock until it can reach the inventory store.
Failure assumptions
- Two regions decrement the same SKU during a partition.
- Both replicas observe a merged note and might each emit a receipt.
- Tombstones from cart removes accumulate for a replica that is offline all day.
Constraints
- Do not put the payment on a PN-Counter.
- Do not introduce two-phase commit across regions to feel safer.
Prompt
A multi-region retail app has likes, a shared shopping note, a cart, and checkout. Checkout must not sell the same unit twice and must not capture a payment twice. Regions partition.
API
Which endpoints merge, and which endpoint reserves stock?
Data
Which Redis mode, if any, holds likes, and which table holds the reservation?
Architecture
Where does the outbox sit relative to the note's publish?
The same product, two planes
Prefer
Merge the draft, commit the invariant
Likes, notes, and carts converge. Usernames, stock, and payments reject all but one winner in a linearizable store.
- Each field has a type or an explicit refusal.
- Publish freezes a version, then emits one event.
- A partition keeps the merge plane up and pauses the commit plane.
Alternative
One mechanism for every write
Either every key waits on a quorum, or every key is a CRDT, including the ones that must not both succeed.
- Keystrokes do not need Raft.
- The last unit of stock does not need a counter.
- Two-phase commit across regions is a third wrong default.
Annotate the field before you pick the cluster
The production failure is a field with no policy, replicated by whatever the product already used.
- 1
Write the invariant in one sentence
Both increments stay, the tag can return, the draft converges, or exactly one of these creates may exist. The sentence picks the plane. - 2
Put merge-shaped fields on a CRDT
Counter, register, observed-remove set, map, or a sequence library. Record the type next to the field so the next reviewer does not swap in last-writer-wins. - 3
Put reject-shaped fields on one commit
Unique index, reservation, payment id. That store is linearizable for the decision. It is often one region, not a geo-CRDT. - 4
Budget the metadata and the lag
Tombstones, delta buffers, and p99 sync delay are the operational objects. A replica that cannot ack is a capacity incident, not a quiet success.
Where these systems actually show up
Decisions
- 1
1 Name the invariant
- next2 Exactly one winner?
- ?
2 Exactly one winner?
- no3 CRDT field or cart
- yes4 Linearizable store
- 3
3 CRDT field or cart
- 4
4 Linearizable store
Lesson map
Production CRDTs & When NOT to Use Them (Riak, Redis CRDT, collab apps, vs
Riak data types, Redis Active-Active, and Yjs or Automerge are production CRDTs. Unique names, non-negative money, and exactly-once effects still belong on a linearizable store.
Architecture. Architecture
Select a node to see why it exists, or an edge to see the protocol, direction, effect, and consequence.
Mermaid export
flowchart TB a["1 Name the invariant"] b["2 Exactly one winner?"] c["3 CRDT field or cart"] d["4 Linearizable store"] a -->|1 Name the invariant| b b -->|no| c b -->|yes| d
The left branch is likes, presence, flags, carts, and drafts. The right branch is usernames, stock, and payments. A real product runs both. The mistake is using the left branch's store for the right branch's rule because it was already multi-region.
Riak data types
Riak shipped counters, sets, maps, and registers as data types, composed so a document could merge field by field instead of storing opaque siblings and hoping the application picked one. The store is the multi-master, hinted-handoff, anti-entropy kind of system. The lesson that stuck: the type has to match the field. A counter field that was a last-writer-wins blob will keep losing increments no matter how careful the gossip is. Map and set churn is the capacity problem. Those are interview stories, not a suggestion to stand up Riak for a new system. The types themselves are the earlier pages.
Redis, three different sentences
| Mode | What a write means | Typical job |
|---|---|---|
| One primary, replicas lag | The primary orders the key | Cache, sessions, queues |
| Active-Active CRDT | Regions merge with a type rule | Geo-local writes that can converge |
| A consensus log somewhere else | One total order | Leadership, a ledger, a unique name |
Active-Active (the CRDB product) gives counters, sets, and hashes a conflict-free rule so two regions can write without electing a single primary for that key. INCR across those regions behaves like the counter on the counters page, not like INCR on one primary. Concurrent DECR can violate "never below zero." You add a reservation, or you stop calling it inventory.
Cache-aside is the other Redis. It sits in front of a database, and the hard part is a stampede of fillers and a key that outlives the row. It does not merge two writers. "We use Redis, so we are eventually consistent" is how those three rows get collapsed. Name the mode.
A collaborative app, as a skeleton
- The draft is a Yjs or Automerge document in the clients, with a server mirror if you need one.
- Transport is a WebSocket or WebRTC provider. Joining the document is an authorization check. The CRDT will not do that check for you.
- Durability is a binary snapshot plus the update log, in a database or object storage. The plain text alone is not enough to merge the next update. That layout is the dissemination page.
- Publish writes one immutable version to a linearizable store and emits the domain event once, through an outbox. Event-driven architecture is that outbox. Do not let every replica that applied the edit send the email.
- Presence is ephemeral. Do not retain it for the life of the document.
A malicious peer who can speak the protocol can craft updates. Sign them, or terminate the provider on a server that checks them, especially if the document is not public.
Flow
- 1
1 CRDT draft converges
- next2 Publish freezes a version
- 2
2 Publish freezes a version
- next3 Ledger or unique index
- 3
3 Ledger or unique index
- next4 Outbox emits one fact
- 4
4 Outbox emits one fact
When a CRDT is the wrong tool
| Need | Use instead | Why the join fails |
|---|---|---|
| One username | A unique index or a linearizable compare-and-set | Both creates succeed, then converge as two members or a later cleanup |
| Money that must not go negative | A ledger transaction, or a reservation | A PN-Counter applies every concurrent decrement |
| Ship or charge once | An idempotency key on one commit, then an outbox | Both replicas can observe the merged state and fire the effect |
| An audit total order | A consensus log | Merge order is not the business order |
| Cross-object integrity | One database transaction, rarely two-phase commit, or a saga | Composed CRDTs are not serializability |
| One writer, low churn | A row version | You paid for metadata and bought nothing |
Raft is the usual total order behind the unique name, the ledger, or the membership of the cluster itself. Two-phase commit is the atomic commit across resource managers you operate, and it blocks while the decision is unknown. It is a poor rescue for a wide-area product that found CRDTs too weak. If both sides are tables in one database, use that database's commit. If payment, inventory, and shipping are separate services that can undo, the workflow is a saga: local commits plus compensations. A CRDT cannot refund a card.
MVCC still matters on the commit plane. It versions the row inside that database. Storing a CRDT blob in an MVCC row gives you local durability of the blob. It does not merge two concurrent JSON values written by two primaries. If you needed that merge, the value type has to be a CRDT. If you needed a reject, MVCC plus a unique index is the reject, and the CRDT should not be in the path.
Hybrids that are the default, not a compromise
- Draft, then publish. Editors converge. Publish freezes a version and is the only write the rest of the company sees.
- Cart, then checkout. An OR-Map holds SKUs and quantities. Inventory reservation is a database or Raft-backed counter with a reject. The set page is the cart type.
- Metrics, then billing. G-Counters and PN-Counters feed a dashboard. The invoice reads a ledger.
- Flags. An observed-remove set can fan a flag out to edges. A kill switch that must turn off everywhere exactly once can still live in a consensus config if the product cannot tolerate a lagging region.
What you operate
- The design doc names a CRDT type, or an explicit "not a CRDT," on every replicated field.
- Sync lag has an SLO. Anti-entropy delay and provider delay, not only process uptime.
- Payload and tombstone counts page someone. The dissemination page is why they grow.
- A partition drill checks two things: replicas agree after the heal, and the money path still rejected the second spend.
- Update streams are authenticated. An open peer-to-peer mesh is not a confidential document store.
- You can snapshot, rebuild, and drop a tombstone epoch on purpose, with the stability rule written down.
Interview Q&A
Design multi-region like counts.
Answer
Use a G-Counter or PN-Counter. Redis Active-Active or an application-level merge both qualify if you can say which. Ship the raw counts to an analytics rollup if the dashboard does not need the replica vectors. Do not reuse the counter as a wallet. A like that arrives twice is a product decision. A dollar that arrives twice is an incident.
Design collaborative notes with a paid export.
Answer
Edit in Yjs or Automerge. Persist snapshots and updates. Export and publish go through a strong store that writes one immutable version. Charge from a ledger keyed by that version, and emit the event once via an outbox. The CRDT's job ended when the draft converged.
Active-Active inventory when one unit remains?
Answer
Do not let two regions decrement a CRDT and both return success. Reserve the unit in one linearizable inventory, or accept oversell and a compensation saga if the business has already chosen that pain. The CRDT cart can still hold the shopper's intent. Checkout is the transition off the CRDT.
Is Cosmos DB last-writer-wins a CRDT?
Answer
The register policy is a CRDT: one value, one timestamp, a silent loser. Custom conflict handling is a procedure you have to read. It is not an observed-remove set, and it is not a sequence. Say "LWW register" if that is the policy you configured.
How do CRDTs sit next to MVCC?
Answer
Put the blob in a row and the database will version that row for local atomicity and snapshots. Two writers in two regions are outside that snapshot. Either they merge with a CRDT rule or one write is last-writer-wins at the row. MVCC explains the anomalies inside one database. It will not merge your document.
Why did Active-Active memory climb after a quiet week?
Answer
Removes and version vectors. The user-visible value shrank while tombstones and unacked deltas remained, often because one region was behind. Look at payload bytes and tombstone counts before you scale the hot key. The fix is the stability and snapshot path, not a bigger counter.
When would you reach for two-phase commit or a saga instead?
Answer
Two-phase commit only if you operate the resource managers, they can prepare, and you accept blocking. That is uncommon on a user-facing multi-region path. A saga fits when checkout really is payment plus inventory plus shipping and each step can commit and later compensate. Neither is a more "correct" CRDT. They are different contracts. Use them when the contract matches, and prefer one database when it already contains both rows.
What do you refuse to replicate to every edge?
Answer
Secrets, and any payload a reader at every replica should not see. CRDT gossip is a fan-out. Authorization does not fall out of the merge. Also refuse a kill switch whose lag would be an outage: that one can live on a consensus config even if the rest of the flag set is an OR-Set.
Pitfalls
- Active-Active as a ledger.
- Cache-aside described as a CRDT, or the reverse.
- Two regions confirming the last unit.
- A publish path that runs on every replica.
- Two-phase commit across the public internet because the CRDT felt weak.
- No tombstone budget, then surprise at the memory graph.
- An open provider on a private document.
Take likes, a shared note, a username, and the last unit of a SKU. Place each on the merge plane or the commit plane. For the note, say where the paid export is recorded. For the SKU, say what a second region must not do during a partition.
Go Deeper
- Redis Active-Active and the application notes for increments that join, and for the commands that do not.
- Yjs, Automerge, and local-first software for the editor skeleton.
- Kleppmann, CRDTs: The Hard Parts, for the production gaps that remain after the type is "correct."
Next: back to the hub.