Replication
Esta página aún no está disponible en tu idioma.
DistributedData replicates state by gossip:
- Every
gossipInterval, each node picks one random peer. - The node sends a full snapshot of every key it currently holds — or as much of it as fits one frame, resuming next tick (see What gets gossiped).
- The peer merges each entry; if its state changed, notify local subscribers.
This means writes propagate eventually — usually within 1-2 gossip rounds (1-2 seconds default). For most workloads this is fine; for cases where it isn’t, see Quorum reads/writes.
The flow
Section titled “The flow”Local writes are immediate. Gossip is for propagation.
Configuration
Section titled “Configuration”const distributedDataOptions = DistributedDataOptions.create().withGossipInterval(1_000);const dd = system.extension(DistributedDataId).start( cluster, distributedDataOptions, // default 1s);The default gossipInterval is 1 second — chosen as a balance
between propagation speed and gossip-bandwidth cost. Lower
intervals → faster convergence + more chatter. Higher → less
chatter, slower convergence.
For typical applications:
| Workload | Interval |
|---|---|
| Latency-sensitive shared state (online presence) | 250-500 ms |
| Default for most apps | 1 s |
| Counters / flags that change rarely | 2-5 s |
| Very large clusters where bandwidth matters | 5-10 s |
The same knob is settable from application.conf, so a deployment can
tune it without a rebuild — the builder still wins where both are set:
actor-ts.distributed-data.gossip-interval = 2sSee Configuration for the rest of the block.
What gets gossiped
Section titled “What gets gossiped”The full state, sliced to fit one frame — each tick,
DistributedData serializes as many keys as fit
max-gossip-bytes and pushes them to one random peer, resuming
next tick where this one stopped. There’s no delta tracking and no
per-peer receipt history, so a store that fits the budget travels
whole on every tick.
This makes gossip volume proportional to total state size, not update rate. A million-entry map costs its full serialized size every time it’s gossiped, whether or not it changed; a small counter costs almost nothing no matter how often it increments.
Slicing is safe because merging is per key, and a key absent from a frame means “no information about it this round” — not “deleted”. CRDT merge is idempotent and commutative, so a key that arrives a tick later, or twice, changes nothing about where the cluster converges; it only changes when.
Why the budget exists
Section titled “Why the budget exists”max-gossip-bytes defaults to 1 MiB and is always clamped down
to remote.max-frame-bytes (16 MiB by default). Without it a
large store produced one frame past the wire cap, and that frame is
rejected by the receiver on its 4-byte length prefix — before a
single key is parsed. The receiving transport answers that by
dropping the connection, which carries heartbeats, membership
gossip and every cross-node tell as well. So the store did not
converge slowly; it did not converge at all, while killing one peer
link per tick.
A single key whose own encoding exceeds the budget cannot be sliced
any further — the unit is one key’s full state. It is skipped, and
the replica warns (rate-limited, naming the key and both sizes) and
increments distributed_data_gossip_skipped_keys_total. That
key does not converge. Split the value, or raise both
max-gossip-bytes and remote.max-frame-bytes.
If your store grows large, raising gossip-interval spreads the
cost over time; lowering max-gossip-bytes spreads one tick’s cost
over more ticks.
Per-peer round selection
Section titled “Per-peer round selection”// Every gossip tick, the node picks ONE random reachable peer.Round-robin gossip would be more predictable but creates synchronized waves; random-peer-per-tick disperses load and typically converges in O(log N) rounds for N peers.
After K rounds, the probability that any specific peer hasn’t
received the update is (1 - 1/N)^K — for N=10 nodes and K=5
rounds, that’s < 60 % chance a peer hasn’t seen it; after K=20
rounds, < 13 %.
In practice, convergence is much faster because gossip is multi-hop: peers re-gossip what they received.
Subscribe to changes
Section titled “Subscribe to changes”const unsubscribe = dd.subscribe<GCounter>('hits', (counter) => { console.log(`hits is now ${counter.value()}`);});
// Later:unsubscribe();Subscribers fire synchronously after every successful merge
that changes the local value (deep-equal check via the CRDT’s
toJSON).
This means:
- Local updates → subscriber fires immediately.
- Remote updates → subscriber fires when gossip arrives + the merge changes the local value.
- Idempotent updates → subscriber does NOT fire (no change).
For dashboard widgets, real-time UI updates, or business logic
that should react to changes anywhere in the cluster, subscribe
is the hook.
When a peer leaves
Section titled “When a peer leaves”When MemberRemoved fires for a cluster member, DistributedData
does nothing — there’s no per-peer state to clean up. Gossip
holds no version vector and no per-peer receipt history; each tick
just picks a random peer from the current up-members, so a
departed node simply stops being a gossip target.
State the peer contributed stays in the local replica, merged in
like any other update. CRDT tombstones (a removed ORSet
element, a null register in an LWWMap) are never garbage
collected — they’re kept permanently so a stale gossip can’t
resurrect a removed key. Over a long-lived store with heavy
churn, tombstone accumulation is the practical cost to watch.
A replica that rejoins under the same identity later (stable pod names, persistent volumes) needs no special handling — it’s just another up-member again, and gossip reconverges its state over the next few ticks.
Bandwidth costs
Section titled “Bandwidth costs”Rough shape for a 10-node cluster, default 1-second gossip:
- Per tick, per node — one message carrying up to
max-gossip-bytesof serialized store goes to a single random peer. Its size is set by total state, not by how much changed since the last tick, so an idle store still pays this each tick; the cost just doesn’t grow with write rate. - Scaling knobs — store size drives the per-message cost,
gossip-intervaldrives how often you pay it, andmax-gossip-bytescaps any single payment. A large store on a fast interval is the expensive combination; a large store with a small budget is the slow-convergence one.
For very-large clusters (50+ nodes), gossip volume can become
significant. The framework’s gossip is anti-entropy style —
all-to-all over time — which scales O(N) per node per round.
Bigger clusters may want larger gossipInterval or a different
gossip topology (not currently configurable; an issue if needed).
Diagnosing slow propagation
Section titled “Diagnosing slow propagation”// Node-A writes:dd.update('hits', ..., (c) => c.increment('a', 1));
// Node-B reads, but the value doesn't reflect:console.log(dd.get('hits')?.value()); // ← shows stale valueIf this surprises you:
- Wait a gossip round (1 s default). Most apparent staleness resolves within 1-2 gossip cycles.
- Check
gossipInterval— if set higher than default, the wait is proportional. - Use
getAsyncwithconsistency: 'majority'for reads that must reflect the latest known state across replicas. - Check cluster membership — if node-B is unreachable from node-A, no gossip is flowing.
Where to next
Section titled “Where to next”- Distributed data overview — the bigger picture.
- Quorum reads/writes — the consistency knobs for stronger guarantees.
- Durable storage — persisting replica state to disk for restart recovery.
- Cluster overview — the membership underneath.
