Zum Inhalt springen
Deutsch

Replikation

DistributedData repliziert Zustand via Gossip:

  • Alle gossipInterval wä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.

node-Bnode-Anode-Bnode-Alokales Update angewendet — sofortin lokalen Zustand mergenper-CRDT-Merge-Funktionfalls sich der lokale Wert geändert hatSubscriber benachrichtigenGossip-Tick — vollständiger Snapshot kommt an

Lokale Writes sind sofort. Gossip ist für die Propagation.

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:

WorkloadIntervall
Latenz-empfindlicher Shared State (Online-Präsenz)250-500 ms
Default für die meisten Apps1 s
Counter / Flags, die sich selten ändern2-5 s
Sehr große Cluster, wo Bandbreite zählt5-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 = 2s

Den Rest des Blocks findest du unter Konfiguration.

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.

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.

// 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.

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 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.

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-bytes an 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-interval treibt, wie oft du sie zahlst, und max-gossip-bytes begrenzt 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).

// 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 Wert

Wenn dich das überrascht:

  1. Warte eine Gossip-Runde (Default 1 s). Die meiste sichtbare Veraltung löst sich in 1-2 Gossip-Zyklen.
  2. Prüfe gossipInterval — wenn höher als der Default gesetzt, ist die Wartezeit proportional.
  3. Nutze getAsync mit consistency: 'majority' für Reads, die den letzten bekannten Stand über die Replikas hinweg reflektieren müssen.
  4. Prüfe Cluster-Membership — wenn node-B von node-A aus unerreichbar ist, fließt kein Gossip.