skip to content

When operating CRDTs at scale in production, what are the concrete costs (metadata growth, tombstone accumulation, causal stability tracking) and in what situations should you avoid choosing a CRDT for a piece of state at all?

level: principalimportance: should knowfreq 35%

answer

  1. tombstone GC needs causal stability
  2. coordination deferred, not eliminated
  3. no invariant enforcement at write time
  4. single-writer or consensus for invariants
  5. Dynamo: cart=CRDT-friendly, payment=not

basics

~20 s

CRDTs make merging automatic, but the bookkeeping (unique tags, deleted-item markers, per-replica counters) keeps growing and someone has to periodically clean it up safely, which is tricky. And CRDTs are the wrong tool whenever you need a strict rule across the whole system at once, like 'account balance can never go negative,' because no replica can enforce that alone without talking to the others first.

solid answer

~50 s

In production, CRDT metadata (tags in an OR-Set, tombstones, per-node vectors) grows monotonically and must eventually be garbage-collected; that's only safe once you know every replica has observed a given tombstone, which requires tracking 'causal stability' - a form of distributed agreement about what's been fully propagated - reintroducing coordination for GC even though writes stayed coordination-free. Left unmanaged, this metadata can dwarf the actual payload. More fundamentally, CRDTs are the wrong fit whenever an operation's correctness depends on a global invariant that spans replicas - e.g. 'don't overdraw this account,' 'don't oversell this inventory slot,' 'enforce a unique username' - because CRDT merges are, by construction, local and coordination-free, so no replica can know at write time whether merging will violate an invariant; those cases need genuine coordination (consensus, single-writer partitioning, or reconciliation with compensating actions) instead.

go deeper

for a junior

Should recognize, at a high level, that 'automatic merging' doesn't mean 'no cleanup needed' and that CRDTs aren't a good fit for things like bank balances.

for a middle

Should be able to say why tombstones/tags can't be deleted immediately and name at least one operational cost of running CRDTs long-term.

for a senior

Should explain causal stability as the mechanism blocking safe GC, and articulate precisely why CRDT merges can't enforce cross-replica invariants at write time.

for a principal

Should be able to design around the limitation - choosing single-writer sharding, consensus, or saga/compensation per field based on its invariant needs - and reason about GC subsystem operational cost as a first-class part of adopting CRDTs at scale, citing the general shopping-cart-vs-payment framing.

## Why the metadata is never safely deletable at write time Every CRDT design — G-Counter's per-node vector, OR-Set's per-add tags and tombstones, PN-Counter's doubled vectors — shares a structural property: none of their internal metadata is ever safely deletable at write time, because deletion is itself a form of coordination the CRDT was built to avoid. A tombstone can only be discarded once every replica that might still hold a stale reference to the tagged add has demonstrably observed the tombstone — otherwise a late-arriving stale replica could resurrect a deleted element, since the OR-Set presence check depends on the tombstone still being there to cancel the corresponding tag. Determining "has every replica observed this" is a form of distributed agreement called **causal stability**, typically tracked via version vectors exchanged periodically among all replicas; it's real coordination work, just deferred from the write path (where CRDTs avoid it) to a background garbage-collection path (where it reappears). In practice this means running CRDTs at scale requires operating a second subsystem — the GC/compaction process — that most teams underestimate when they first adopt CRDTs specifically to escape coordination. ## No free lunch This cost exists precisely because it's the flip side of the coordination-free write guarantee: a system that lets every replica accept writes locally, instantly, without contacting any other replica, cannot simultaneously know in real time which of its own past writes are safe to forget, because "safe to forget" requires knowing what every other replica has (or hasn't) seen — information that, by construction, wasn't exchanged at write time. There is no free lunch here: either you pay coordination cost per write (traditional replicated systems, consensus-based systems), or you defer it to background reconciliation/GC (CRDTs), but a genuinely coordination-free, garbage-free system for arbitrary concurrent state does not exist. ## The operational trade-off The practical trade-off in production is between: - **unbounded metadata growth** (skip GC entirely, and merge time/storage/network bandwidth for the CRDT slowly degrades as tombstones and tags accumulate — fine for low-churn state, a real problem for high-churn state like presence indicators or frequently-retagged items); - **versus running a periodic causal-stability protocol** (bounded metadata, at the cost of building and operating that background subsystem correctly, including handling the case where a replica is permanently gone and its "pending acknowledgment" would otherwise block GC forever — most real implementations need an explicit replica-removal/eviction procedure for this). Systems that ship CRDTs as a first-class primitive (Riak, Akka Distributed Data) bake some form of this GC or pruning into the platform; teams that hand-roll a CRDT on top of a generic datastore often don't, and pay for it later as tombstone bloat degrades performance. ## The most consequential failure mode The most consequential failure mode, though, isn't operational overhead — it's choosing a CRDT for state that carries a **cross-replica invariant**. CRDTs guarantee that merges are well-defined and always succeed, but "always succeeds" is exactly the property that makes them unable to reject an update based on global state: a PN-Counter merge can't refuse to let the value go negative, an OR-Set add can't check "is this username already taken system-wide" before accepting, because checking either would require contacting other replicas at write time — the very coordination CRDTs are designed to avoid. Teams that reach for a PN-Counter for account balance, or an OR-Set for a "claimed usernames" set, discover the invariant violation in production: - two partitioned replicas both let a balance go to -50 because each only saw its own writes; - two users both successfully claim "admin" as their username in different regions during a partition. Reconciling that after the fact (deciding which claim "really" wins, refunding an overdraft) is exactly the manual, application-specific conflict resolution CRDTs were supposed to eliminate — just moved to a worse place, after the invariant has already been visibly violated to users. ## The rule of thumb The rule of thumb: **use a CRDT** when the operation's result is genuinely well-defined independent of what else concurrently happened elsewhere (accumulating a count, unioning a set of tags, converging on a "last write" for a display field) and availability-under-partition matters more than a real-time global constraint. **Don't use a CRDT** whenever the correctness of an update genuinely depends on the current state at every other replica, such as account balances that can't go negative, oversell-sensitive inventory counts, or global uniqueness constraints like usernames or email addresses. Use, instead: 1. **single-writer sharding** (route all writes for a given key to one owning replica/partition); 2. **a consensus protocol** (Raft/Paxos-backed strongly consistent store) for the invariant-bearing field; 3. **a saga/compensating-transaction pattern** (accept the invariant might be briefly violated, detect it, and issue a corrective action afterward). Amazon's original **Dynamo** paper, which inspired much of this line of work including Riak, is explicit about exactly this trade: shopping-cart "add item" is CRDT-friendly (an OR-Set of cart items merges fine, worst case a removed item reappears and gets silently re-removed on checkout), but payment capture is not, and is handled by a strongly consistent path instead.

  • If GC requires coordination anyway, doesn't that defeat the purpose of using a CRDT?
    Not entirely - the coordination for GC happens off the write path, in the background, asynchronously, and its failure or delay only affects storage/bandwidth efficiency, not write availability; writes remain fully available and coordination-free even if GC is temporarily stalled or skipped, which is a meaningfully different (and usually acceptable) failure mode compared to writes themselves requiring coordination.
  • Could you enforce an invariant like non-negative balance by rejecting reads that show a negative value instead of preventing the write?
    That papers over the problem rather than solving it - by the time a negative value is visible, the invariant has already been violated at the data layer; you'd need application logic to detect and correct it after the fact (a compensating transaction, like reversing the overdraft and notifying the user), which is a valid pattern but a fundamentally different, after-the-fact reconciliation approach rather than genuine invariant enforcement.
  • Are there CRDTs that support bounded counters (e.g. can't exceed a cap)?
    There's research on bounded counter CRDTs that pre-allocate a fixed 'budget' of allowed increments/decrements to each replica upfront and let replicas trade budget with each other, which can enforce a bound without per-write coordination, but it requires knowing the bound and replica set in advance and adds real complexity - it's a specialized technique, not a drop-in replacement for PN-Counter, and most production systems facing this problem choose a coordinated approach instead.

It's like a shared expense-splitting spreadsheet that lets everyone edit their own row offline and merges automatically when everyone's back online - great for tallying who owes what, but useless for enforcing 'the group's account can never go below zero,' because no single editor, working offline, can know in the moment whether their edit will push the shared total negative once everyone else's edits are combined.

saying these in an interview costs you the question

  • Assumes CRDTs need no garbage collection at all
  • Proposes a CRDT for account balance or inventory without caveats
  • Doesn't recognize GC as reintroducing coordination
  • Thinks metadata growth is negligible regardless of churn rate
  • Can't name an alternative (single-writer, consensus, saga) for invariant-bearing state

context