Replikation
DistributedData repliziert Zustand via Gossip:
- Alle
gossipIntervalMswählt jeder Node einen zufälligen Peer. - Der Node sendet einen vollständigen Snapshot jedes Keys, den er aktuell hält.
- Der Peer merged jeden Eintrag; wenn sich sein Zustand geändert hat, benachrichtigt er die lokalen Subscriber.
Das bedeutet, Writes propagieren eventually — typischerweise in 1-2 Gossip-Runden (1-2 Sekunden default). Für die meisten Workloads ist das ausreichend; für Fälle, in denen es das nicht ist, siehe Quorum Reads/Writes.
Der Ablauf
Abschnitt betitelt „Der Ablauf“Lokale Writes sind sofort. Gossip ist für die Propagation.
Konfiguration
Abschnitt betitelt „Konfiguration“const distributedDataOptions = DistributedDataOptions.create().withGossipInterval(1_000);const dd = system.extension(DistributedDataId).start( cluster, distributedDataOptions, // default 1s);Der Default für gossipIntervalMs ist 1 Sekunde — gewählt als
Balance zwischen Propagation-Geschwindigkeit und Gossip-Bandbreite.
Kleinere Intervalle → schnellere Konvergenz + mehr Traffic.
Größere → weniger Traffic, langsamere Konvergenz.
Für typische Anwendungen:
| Workload | Intervall |
|---|---|
| Latenz-empfindlicher Shared State (Online-Präsenz) | 250-500 ms |
| Default für die meisten Apps | 1 s |
| Counter / Flags, die sich selten ändern | 2-5 s |
| Sehr große Cluster, wo Bandbreite zählt | 5-10 s |
Derselbe Knopf ist auch aus der application.conf setzbar, ein Deployment
kann ihn also ohne Rebuild tunen — wo beides gesetzt ist, gewinnt weiterhin
der Builder:
actor-ts.distributed-data.gossip-interval = 2sDen Rest des Blocks findest du unter Konfiguration.
Was gegossipt wird
Abschnitt betitelt „Was gegossipt wird“Der vollständige State — bei jedem Tick serialisiert DistributedData jeden Key seiner lokalen Replika und pusht den gesamten Snapshot an einen zufälligen Peer. Es gibt kein Delta-Tracking und keine Empfangs-Historie pro Peer; jedes Gossip trägt die komplette Menge.
So ist das Gossip-Volumen proportional zur Gesamt-State-Größe, nicht zur Update-Rate. Eine Map mit einer Million Einträgen kostet jedes Mal, wenn sie gegossipt wird, ihre volle serialisierte Größe — egal ob sie sich geändert hat; ein kleiner Counter kostet fast nichts, egal wie oft er inkrementiert wird.
Das ist bewusst simpel — kein Digest, kein Delta — was billig zu
implementieren und für die kleinen bis mittleren Stores, für die
DistributedData gedacht ist, gut genug ist. Wenn dein Store groß
wird, erhöhe gossipIntervalMs, um die Kosten über die Zeit zu
verteilen.
Auswahl des Peers pro Runde
Abschnitt betitelt „Auswahl des Peers pro Runde“// Bei jedem Gossip-Tick wählt der Node EINEN zufälligen erreichbaren Peer.Round-Robin-Gossip wäre vorhersagbarer, erzeugt aber synchronisierte Wellen; ein zufälliger Peer pro Tick verteilt die Last und konvergiert typischerweise in O(log N) Runden bei N Peers.
Nach K Runden ist die Wahrscheinlichkeit, dass ein bestimmter
Peer das Update nicht erhalten hat, (1 - 1/N)^K — für N=10
Nodes und K=5 Runden also < 60 % Wahrscheinlichkeit, dass ein
Peer es nicht gesehen hat; nach K=20 Runden < 13 %.
In der Praxis ist die Konvergenz viel schneller, weil Gossip multi-hop ist: Peers regossipen, was sie erhalten haben.
Auf Änderungen subscriben
Abschnitt betitelt „Auf Änderungen subscriben“const unsubscribe = dd.subscribe<GCounter>('hits', (counter) => { console.log(`hits liegt jetzt bei ${counter.value()}`);});
// Später:unsubscribe();Subscriber feuern synchron nach jedem erfolgreichen Merge, der
den lokalen Wert ändert (Deep-Equal-Check über das toJSON des
CRDT).
Das heißt:
- Lokale Updates → Subscriber feuert sofort.
- Remote Updates → Subscriber feuert, wenn Gossip ankommt und der Merge den lokalen Wert ändert.
- Idempotente Updates → Subscriber feuert NICHT (keine Änderung).
Für Dashboard-Widgets, Echtzeit-UI-Updates oder Business-Logik,
die auf Änderungen irgendwo im Cluster reagieren soll, ist
subscribe der Hook.
Wenn ein Peer den Cluster verlässt
Abschnitt betitelt „Wenn ein Peer den Cluster verlässt“Wenn MemberRemoved für ein Cluster-Mitglied feuert, tut
DistributedData nichts — es gibt keinen Zustand pro Peer, der
aufzuräumen wäre. Gossip hält keinen Version Vector und keine
Empfangs-Historie pro Peer; jeder Tick wählt einfach einen
zufälligen Peer aus den aktuellen Up-Members, sodass ein
ausgeschiedener Node schlicht kein Gossip-Ziel mehr ist.
Der Zustand, den der Peer beigetragen hat, bleibt in der lokalen
Replika, eingemerged wie jedes andere Update. CRDT-Tombstones
(ein entferntes ORSet-Element, ein null-Register in einer
LWWMap) werden nie garbage-collected — sie bleiben dauerhaft
erhalten, damit ein veraltetes Gossip einen entfernten Key nicht
wiederauferstehen lassen kann. Über einen langlebigen Store mit
starkem Churn ist die Tombstone-Akkumulation der praktische
Kostenpunkt, den man im Auge behalten sollte.
Eine Replika, die später unter derselben Identität wieder beitritt (stabile Pod-Namen, persistente Volumes), braucht keine Sonderbehandlung — sie ist einfach wieder ein weiterer Up-Member, und Gossip lässt ihren Zustand über die nächsten paar Ticks erneut konvergieren.
Bandbreitenkosten
Abschnitt betitelt „Bandbreitenkosten“Grobe Zahlen für einen 10-Node-Cluster mit Default-Gossip von 1 Sekunde:
- Pro Tick, pro Node — eine Nachricht, die den vollständig serialisierten Store trägt, geht an einen einzigen zufälligen Peer. Ihre Größe wird durch den Gesamt-State bestimmt, nicht dadurch, wie viel sich seit dem letzten Tick geändert hat, sodass ein Idle-Store das trotzdem bei jedem Tick zahlt; die Kosten wachsen nur nicht mit der Write-Rate.
- Skalierungs-Stellschrauben — die Store-Größe treibt die
Kosten pro Nachricht und
gossipIntervalMstreibt, wie oft du sie zahlst. Ein großer Store bei einem schnellen Intervall ist die teure Kombination.
Für sehr große Cluster (50+ Nodes) kann das Gossip-Volumen
relevant werden. Das Gossip des Frameworks ist im
Anti-Entropy-Stil — all-to-all über die Zeit — was O(N) pro
Node pro Runde skaliert. Größere Cluster wollen vielleicht ein
größeres gossipIntervalMs oder eine andere Gossip-Topologie
(aktuell nicht konfigurierbar; bei Bedarf ein Issue).
Langsame Propagation diagnostizieren
Abschnitt betitelt „Langsame Propagation diagnostizieren“// Node-A schreibt:dd.update('hits', ..., (c) => c.increment('a', 1));
// Node-B liest, aber der Wert spiegelt es nicht wider:console.log(dd.get('hits')?.value()); // ← zeigt veralteten WertWenn dich das überrascht:
- Warte eine Gossip-Runde (Default 1 s). Die meiste sichtbare Veraltung löst sich in 1-2 Gossip-Zyklen.
- Prüfe
gossipIntervalMs— wenn höher als der Default gesetzt, ist die Wartezeit proportional. - Nutze
getAsyncmitconsistency: 'majority'für Reads, die den letzten bekannten Stand über die Replikas hinweg reflektieren müssen. - Prüfe Cluster-Membership — wenn node-B von node-A aus unerreichbar ist, fließt kein Gossip.
Wohin als Nächstes
Abschnitt betitelt „Wohin als Nächstes“- Distributed Data im Überblick — das Gesamtbild.
- Quorum Reads/Writes — die Konsistenz-Stellschrauben für stärkere Garantien.
- Durable Storage — Replika-State für Restart-Recovery auf Disk persistieren.
- Cluster im Überblick — das Membership-Modell darunter.
