Distributed systems
Part 5 of 6 · CRDTsState-based vs Op-based vs Delta-CRDTs & Compaction
The same abstract type can ship as full state, as operations, or as deltas. The channel, the fresh replica, and the tombstone pile follow from that choice.
- 1Gist
- 2Maps
- 3Q&A
- 4Sandbox
Voice readout needs Web Speech Synthesis in this browser.
Question ladder
L1
What does a state-based replica send?
Answer
A mergeable state. The receiver joins it. Resending is safe because join is idempotent. Order does not matter because join is commutative and associative.
L2
What does an op-based replica require from the channel?
Answer
Reliable causal delivery. Concurrent ops must commute. If the op is not idempotent, duplicates must be filtered by a unique op id.
L3
What is a delta?
Answer
A join-fragment since the last sync. Peers join it the way they would join full state. Acknowledgements prune the buffer of fragments the peer has already joined.
L4
Can you run an op-based CRDT over UDP with no acks?
Answer
Not under the op-based assumptions. A lost op never arrives, and the replica is missing an update rather than merely stale. Prefer state or deltas on a lossy channel, unless every op is itself an idempotent join.
L5
What breaks if merge is not idempotent?
Answer
A retransmit applies twice. Counters inflate or replicas that saw the same messages disagree. That function is not a CRDT merge.
L6
When may a tombstone be deleted?
Answer
When it is causally stable: every replica in the membership has observed it. Then no legal lagging message still depends on that evidence.
L7
How is this different from Raft log compaction?
Answer
Both snapshot and truncate a log. A Raft snapshot preserves a linearizable state machine. A CRDT snapshot preserves a mergeable state. The consistency goal is not the same, so one snapshot does not replace the other.
Failure modes
Full state on a large set
Fifty thousand elements go out every round because ten of them changed. The cluster spends its budget on bytes that were already known.
Ops on a lossy link
A dropped increment never joins. The sender has moved on. Anti-entropy cannot reconstruct an op that was not retained.
A delta buffer that prunes early
The sender forgets a fragment the peer never acknowledged. The peer's join is missing a dot, and later full states are not in the protocol to repair it.
Tombstone GC by age
A replica offline longer than the TTL syncs an old add. The tombstone is gone. The element returns.
Misconceptions
Delta-CRDTs are op-based because the messages are small.
Deltas are lattice fragments. Join stays idempotent. The channel can still be lossy if you are willing to resend the unacked fragment.
Causal delivery is a total order.
Causal delivery respects happens-before. Concurrent ops still need to commute. A Raft log is the total order, and it is a different protocol.
Anti-entropy is cache invalidation.
Anti-entropy compares replicas and repairs divergence with a merge. Cache-aside decides whether a key is fresh against a database. Those are different bugs.
Interviewer traps
We will just resync the whole database, so the mode does not matter.
Say who pays for that resync, how often, and what a blank replica does. Then pick state, op, or delta on purpose.
CRDT compaction is Raft snapshotting.
The mechanics rhyme. The guarantee does not. Point at the Raft hub for the linearizable log and stay on causal stability here.
Design scenario
Same prompt for every reader.
Requirements
A fresh replica can catch up. A dropped packet must not lose an update forever. Tombstones must not grow without a bound, and must not disappear while a region is offline.
Traffic / scale
Ten updates a minute in steady state, with rare bulk imports of thousands of ids.
Latency
Interactive reads use the local replica. Cross-region convergence within a minute is enough.
Consistency
Strong eventual consistency after the same updates are delivered. Not linearizability.
Availability
Four regions keep accepting writes while one is offline. The offline region merges when it returns, including deletes it missed.
Failure assumptions
- Gossip drops packets.
- A replica is empty and must catch up.
- A region returns after 24 hours with a stale add for a deleted id.
Constraints
- Do not gossip the full set on every tick once you have measured it.
- Do not delete tombstones on a 24-hour timer.
Prompt
An observed-remove set of about 50,000 ids is replicated across five regions. A typical minute changes ten elements. One region can sit offline for a day. Deletes are common.
API
What does a replica send on a steady-state tick, and what does it send to an empty peer?
Data
Which bytes are the live set, which are tombstones, and which are the ack vector?
Architecture
Where is the stability oracle that allows garbage collection?
Ten elements changed in a large set
Prefer
Ship the join of what changed
A delta or an op log sends the ten updates. A snapshot still exists so a blank replica does not replay from the dawn of the set.
- The steady-state message tracks churn, not cardinality.
- A lost delta is resent until it is acknowledged.
- Tombstones stay until the offline region has seen them.
Alternative
Ship the whole set every round
It converges. It also sends a megabyte so the receiver can learn about ten ids.
- Simple, and the right prototype.
- The wrong steady state once you have measured it.
- It still does not let you drop tombstones early.
Pick the channel, then the encoding
A small message is not automatically an op. Ask whether a duplicate is safe.
- 1
Write down what the network may do
Drops, duplicates, and reordering want an idempotent join. A reliable causal channel can carry operations that are not safe to apply twice. - 2
Send state, a delta, or an op
Full state is the join of everything. A delta is the join since an acknowledgement. An op is one mutation plus the causal dependencies it needs. - 3
Define catch-up for an empty replica
State-based catch-up is 'ship the state.' Op-based catch-up is a snapshot plus the log, or a fallback state transfer. Do not leave this as a TODO. - 4
Compact only stable evidence
Snapshot, truncate the log before that snapshot, and drop tombstones whose dots every member has observed.
Three ways to ship one type
An OR-Set does not care, at the semantic level, whether the bits on the wire were a full state, a delta, or a stream of adds and removes. The rest of the system cares. Bandwidth, what a crash loses, and what a brand-new replica must download all change.
Decisions
- 1
1 Local update
- next2 What do you ship?
- ?
2 What do you ship?
- full or delta3 Idempotent join
- operation4 Causal apply
- 3
3 Idempotent join
- next5 Same updates, same state
- 4
4 Causal apply
- next5 Same updates, same state
- 5
5 Same updates, same state
Lesson map
State-based vs Op-based vs Delta-CRDTs & Compaction
The same abstract type can ship as full state, as operations, or as deltas. The channel, the fresh replica, and the tombstone pile follow from that choice.
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 Local update"] b["2 What do you ship?"] c["3 Idempotent join"] d["4 Causal apply"] a -->|1 Local update to 2 What do you ship?| b b -->|full or delta| c b -->|operation| d
| Mode | Message | Channel | Empty replica | Complexity |
|---|---|---|---|---|
| Full state | The object | Loss and reorder are fine | Ship the state | Lowest |
| Operation | One mutation | Reliable causal delivery | Snapshot plus ops | Higher |
| Delta | The join since the last ack | Loss is fine if you resend | Snapshot plus the missing deltas, or the deltas since genesis | Highest |
| Snapshot plus log | Bounded | Mixed on purpose | The usual catch-up | What production ships |
State-based
The sender transmits its state. The receiver replaces its state with merge(local, remote), the least upper bound. Because that function is commutative, associative, and idempotent, the sender can repeat itself, the packets can arrive out of order, and a duplicate does not double the value.
The cost is the size of the object. An OR-Set with 50,000 elements, two dots each, and 16 bytes per dot is about 1.6 MB before headers. Five peers gossiping that every 30 seconds move about 8 MB per round even when ten elements changed. That number is why "just ship the state" stops being a design and becomes an incident once the set is real.
Op-based
The sender transmits Inc("A") or Add(dot, element). The receiver applies the op when its causal dependencies are already in the local set. Concurrent ops must commute. Ops that are ordered by happens-before can rely on that order and need not commute with each other.
The channel does the hard part: no loss, no reorder against happens-before. If the op is not idempotent, a duplicate is a second increment unless you filter by op id. A new replica cannot apply "the latest op" and be correct. It needs a snapshot and the ops after that snapshot, or a one-time state transfer. UDP without acknowledgements does not meet the delivery assumption. If you need that network, you wanted a state-based or delta design.
Deltas
A delta is a fragment of the same join-semilattice, not a different algebra. Joining deltas is associative, commutative, and idempotent, just as joining full counters is. Each peer remembers what it has sent and what the neighbor has acknowledged, and it drops fragments the neighbor has joined. That buffer is the whole trick. Ack too early and the neighbor misses a join forever, unless some other anti-entropy path ships a full state. Ack too late and the buffer is the full state with extra bookkeeping.
Deltas are how you approach "bytes proportional to the change" without giving up a lossy network. They are also where implementation bugs hide, because the failure looks like a slow divergence rather than an exception.
Decisions
- 1
1 Snapshot the state
- next2 Truncate the covered log
- 2
2 Truncate the covered log
- next3 Tombstone stable?
- ?
3 Tombstone stable?
- yes4 Drop that evidence
- no5 Keep the tombstone
- 4
4 Drop that evidence
- 5
5 Keep the tombstone
A worked size comparison
Use the 1.6 MB set from above.
- Full-state gossip to five peers is on the order of 8 MB a round, plus headers, for ten changed elements.
- Ten add or remove ops at about 50 bytes each are about 500 bytes, plus whatever causal metadata the op carries.
- A delta of the changed dots is typically kilobytes when churn is low, and it spikes when someone imports a bulk file. That spike is the thing to measure, not the average.
Production systems default to a snapshot plus deltas or ops because of this arithmetic. The metric is bytes per sync next to the churn rate, not a cluster-wide "replication looks fine."
Compaction
Removes keep evidence: tombstones, removed dots, deleted sequence nodes. Without a plan, memory and payloads grow with churn, and the mobile client is the first casualty.
Causal stability means every current member has observed the remove. Only then is the evidence redundant. Computing that needs a membership and a summary of acknowledgements. In a pure peer-to-peer mesh with no agreed membership, the proof is much harder. Products use epochs, periodic snapshots that define a new baseline, or a server that is allowed to say who the members are.
A practical recipe:
- Snapshot the CRDT in its binary form on a period or on a size threshold.
- Truncate the op or delta log that the snapshot already includes.
- Run tombstone collection only when the stability oracle says the dot is known everywhere you still accept sync from.
- Compress what you still send. Yjs update encoding and Automerge's columnar format exist for this.
- Alert on payload bytes, tombstone count, and replica skew. A count of "merges succeeded" will stay green while the documents become unloadable.
Anti-entropy is not a cache
Riak-style stores combine gossip or hinted handoff, read repair, and active anti-entropy with Merkle trees. The conflict resolver is the CRDT merge, which is why they stopped handing the application a pile of opaque siblings.
That machinery is not Redis cache-aside. Cache-aside fills a key from an origin and worries about a stampede of fillers. It does not join two writers. Active-Active CRDT types are a different Redis product, covered on the next page. If a design review says "anti-entropy" and then describes GET and SET, ask which replica is the origin of truth.
Raft compaction is the other rhyme to refuse. A Raft snapshot lets a follower skip a prefix of a linearizable log. A CRDT snapshot lets a peer skip a prefix of a mergeable history. Both truncate. They do not preserve the same thing. Raft owns the first guarantee.
How you test it
Generate random schedules. Drop messages. Reorder them. Duplicate them. After you flush the channel, every replica that was supposed to receive the updates has the same value. Separately, assert a byte ceiling after garbage collection once you have simulated the acknowledgements. A unit test that merges two counters once will not catch a delta buffer that forgets a fragment.
Sandbox
The Python join is the G-Counter delta. It is the same max-per-slot function as the full state, which is why a duplicate delta is harmless. The TypeScript check is the op-based gate: do not apply an operation until every causal dependency is already in the seen set.
ProblemJoin three partial counter maps. Confirm the join is associative, commutative, and idempotent, and that a stale slot cannot lower a higher one.
ExpectedThe joined slots are A:3, B:2, and C:1, which sum to 6. Joining a fragment with itself changes nothing.
Edge cases
- A missing key is zero.
- Order of joins does not matter.
- Test: associative
join(join(d1, d2), d3) == join(d1, join(d2, d3)) - Test: commutative
join(d1, d2) == join(d2, d1) - Test: idempotent
join(full, full) == full - Test: sum of the joined slots
sum(full.values()) == 6
Press Run. Snippets must be self-contained — no network, files, or native modules.
ProblemAn add that depends on op-1 must wait until op-1 is in the seen set. A duplicate id must not be applied twice if you treat the id as already seen.
ExpectedThe dependent op is blocked, then allowed after op-1 is marked seen. The duplicate is refused.
- Test: missing dep blocks
canApply(addMilk, seen) === false - Test: dep present allows
canApply(addMilk, withDep) === true - Test: duplicate id is already seen
seen.has(addMilk.id) === false && withDup.has(addMilk.id) === true
Press Run. Snippets must be self-contained — no network, files, or native modules.
Interview Q&A
Can I run an op-based CRDT over UDP without acknowledgements?
Answer
Not if you still want the op-based guarantee. A lost op is a lost update. You can resend ops yourself until they are acked, at which point you have built a reliability layer, or you can send idempotent state or deltas and let the join absorb the duplicates. "The op happens to be idempotent" is a special case you should be ready to prove.
Why deltas instead of always shipping state?
Answer
Because full-state gossip is proportional to the object, and the object is usually much larger than the last minute of edits. Deltas aim at the size of the change. They still join, so they tolerate a resend. The price is the acknowledgement buffer.
What breaks if merge is not idempotent?
Answer
Retries. Gossip sends the same payload twice, the receiver adds it twice, and two replicas that both received the same logical update now disagree or have inflated. The function was a merge in the English sense and not a join.
How do you test dissemination rather than the datatype?
Answer
Property tests: random ops, random drops, random reorders, random duplicates. After the channel is flushed, replicas that should have the updates agree. Add a bound on bytes after a simulated garbage collection. A single hand-merged example will not catch a prune bug.
How does this relate to Raft compaction?
Answer
You snapshot and you throw away a prefix. In Raft the prefix is a linearizable log and the snapshot must match what a quorum committed. In a CRDT the prefix is merge evidence and the snapshot must be a state peers can still join. Same shape, different promise. Do not answer one with the other's proof.
What is causal stability in one sentence?
Answer
Every replica you still sync with has observed this remove, so no future message should depend on its tombstone. Membership is part of the sentence. A tombstone older than a day is not the sentence.
What does a blank replica need in each mode?
Answer
State-based: the current state. Delta-based: a snapshot plus every delta the snapshot does not cover, or the full delta history if you never snapshot. Op-based: a snapshot plus the ops after it. "Replay the ops we still have in memory" loses the replica that joins after a restart.
Why mention cache-aside on a dissemination page?
Answer
Because both get called replication. Cache-aside invalidates or refills a key in front of an origin. CRDT anti-entropy merges peers that are all origins for their local writes. Redis cache-aside is the first problem. This page is the second.
Pitfalls
- Sending operations over a channel that drops them.
- Pruning a delta the peer has not acknowledged.
- Garbage-collecting tombstones by age.
- Gossiping full state after the set no longer fits the budget, because the prototype did.
- Calling a Merkle repair pass a cache stampede, or the reverse.
- Treating a CRDT snapshot as a Raft snapshot.
Take 50,000 elements, two dots, 16 bytes. Compute a full-state round to five peers. Then price ten small ops. Say which one you would alert on, and what has to be true before you delete the tombstones from the deletes in that round.
Go Deeper
- Almeida, Shoker, and Baquero, Delta State Replicated Data Types.
- The state-based and op-based sections of Shapiro et al..
- Yjs update encoding, as a concrete snapshot-plus-delta format.