Replikation
DistributedData repliziert Zustand via Gossip:
- Alle
gossipIntervalwählt jeder Node einen zufälligen Peer. - Der Node sendet einen vollständigen Snapshot jedes Keys, den er aktuell hält — oder so viel davon, wie in ein Frame passt, und setzt beim nächsten Tick fort (siehe Was gegossipt wird).
- 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 gossipInterval 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, auf ein Frame zugeschnitten — bei jedem
Tick serialisiert DistributedData so viele Keys, wie in
max-gossip-bytes passen, und pusht sie an einen zufälligen Peer;
beim nächsten Tick geht es dort weiter, wo dieser aufgehört hat. Es
gibt kein Delta-Tracking und keine Empfangs-Historie pro Peer — ein
Store, der ins Budget passt, reist also bei jedem Tick vollständig.
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 Zuschneiden ist sicher, weil pro Key gemerged wird und ein Key, der in einem Frame fehlt, “keine Information über ihn in dieser Runde” bedeutet — nicht “gelöscht”. CRDT-Merge ist idempotent und kommutativ: Ein Key, der einen Tick später oder zweimal ankommt, ändert nichts daran, wohin der Cluster konvergiert, nur wann.
Warum es das Budget gibt
Abschnitt betitelt „Warum es das Budget gibt“max-gossip-bytes hat den Default 1 MiB und wird immer auf
remote.max-frame-bytes heruntergeschnitten (Default 16 MiB).
Ohne das Budget erzeugte ein großer Store ein Frame jenseits des
Wire-Caps — und so ein Frame wird beim Empfänger schon an seinem
4-Byte-Längenpräfix abgelehnt, bevor ein einziger Key geparst
wird. Der empfangende Transport verwirft daraufhin die Verbindung,
und über die laufen auch Heartbeats, Membership-Gossip und jedes
tell über Knotengrenzen. Der Store konvergierte also nicht
langsam, sondern überhaupt nicht — und riss pro Tick eine
Peer-Verbindung mit.
Ein einzelner Key, dessen eigene Kodierung das Budget
überschreitet, lässt sich nicht weiter aufteilen — die Einheit ist
der vollständige State eines Keys. Er wird übersprungen, die
Replika warnt (ratenbegrenzt, mit Key-Name und beiden Größen) und
zählt distributed_data_gossip_skipped_keys_total hoch. Dieser
Key konvergiert nicht. Teile den Wert auf, oder erhöhe
max-gossip-bytes und remote.max-frame-bytes.
Wenn dein Store groß wird, verteilt ein höheres gossip-interval
die Kosten über die Zeit; ein kleineres max-gossip-bytes verteilt
die Kosten eines Ticks über mehr Ticks.
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 mit bis zu
max-gossip-bytesan serialisiertem Store 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,
gossip-intervaltreibt, wie oft du sie zahlst, undmax-gossip-bytesbegrenzt jede einzelne Zahlung. Ein großer Store bei einem schnellen Intervall ist die teure Kombination; ein großer Store bei einem kleinen Budget die langsam konvergierende.
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 gossipInterval 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
gossipInterval— 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.
