Zum Inhalt springen
Deutsch

Replikation

DistributedData repliziert Zustand via Gossip:

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

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

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

// 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, 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 gossipIntervalMs treibt, 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).

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