Rebalance
Wenn sich die Mitgliedschaft des Clusters ändert — ein Node tritt bei oder geht — müssen Shards wandern, um die Last verteilt zu halten. Der Koordinator treibt den Prozess: er wählt Shards aus, die umziehen, sagt der Quellregion, sie soll übergeben, und wartet auf Bestätigung, bevor er neu allokiert.
Das ist absichtlich konservativ — gepufferte Nachrichten
warten auf HandOffComplete, bevor sie an den neuen Besitzer
weitergeleitet werden, damit nichts an einer halb-verschobenen
Entity vorbei rast.
Was Rebalance auslöst
Abschnitt betitelt „Was Rebalance auslöst“Drei Pfade:
- Mitgliedschaftsgetrieben —
MemberUp(ein neuer Node),MemberRemoved(ein gehender Node) oder jeder Cluster-Übergang, der die Kandidatenmenge ändert. Der Koordinator führt die Allokationsstrategie für jeden besessenen Shard erneut aus. - Strategiegetrieben — alle
rebalanceIntervalMs(Default 2 s) fragt der Koordinator die Allokationsstrategie nach ihren Rebalance-Empfehlungen.LeastShardAllocationStrategygibt Shards zum Drainen von beschäftigten Nodes zurück;HashAllocationStrategygibt Shards zurück, deren Hash-Ziel gewandert ist. - Eine Region, die stoppt — eine Region meldet sich beim Koordinator ab, während sie herunterfährt, und ihre Shards werden sofort neu verteilt. Das deckt eine Region ab, die auf einem Node gestoppt wird, der im Cluster bleibt — dort feuert überhaupt kein Membership-Event.
Die ersten beiden speisen den Handoff-Mechanismus unten. Der dritte nicht: Die Region ist bereits weg, es gibt also nichts zu übergeben — der Koordinator streicht sie schlicht aus der Registry und alloziert neu.
Ohne diesen dritten Pfad ließ eine gestoppte Region ihre Shards verwaist zurück: Der Koordinator nannte sie weiterhin als deren Zuhause, Sender cachten dieses Zuhause und stellten in einen gestoppten Actor zu, und die Nachrichten wurden zu Dead Letters. Es erholte sich auch nichts — die Kandidatenmenge wird aus der Registry abgeleitet, ohne Liveness-Prüfung, der Rebalance-Tick sah also einen ausgeglichenen Cluster und verschob nichts. Das blieb so, bis der Node den Cluster ganz verließ.
Die Meldung ist Best-Effort und idempotent: Sie entfällt, wenn die
Region nie einen Koordinator gefunden hat, ein Transportfehler auf dem
Weg nach unten wird geloggt und geschluckt, und der MemberRemoved-
Pfad, der einem kompletten Node-Shutdown folgt, findet die Region
bereits weg und tut nichts.
Die Handoff-Sequenz
Abschnitt betitelt „Die Handoff-Sequenz“Wenn der Koordinator entscheidet, dass Shard X von Node A zu
Node B wandern soll:
HandOff(X)wird an die Region von Node A geschickt.- Die Region von Node A markiert Shard
Xalshanding-off. Sie empfängt weiterhin Nachrichten für Entities inX, aber puffert sie, statt sie an Entity-Actors weiterzuleiten. - Jede Entity in
Xerhält das Stop-Signal des Frameworks. Sie führenpostStopaus, persistieren etwaigen Final-State und sterben. - Sobald alle Entities in
Xweg sind, schickt Node AHandOffComplete(X)an den Koordinator. - Der Koordinator führt
allocate(X, ...)aus, um ein neues Zuhause zu wählen. Nehmen wir an, er wählt Node B. - Der Koordinator veröffentlicht den neuen Besitzer per Gossip. Gepufferte Nachrichten auf Node A beginnen, an die Region von Node B weiterzuleiten.
- Die Region von Node B empfängt Nachrichten für Entities in
Xund spawnt Entities bei Bedarf (genau wie bei der Erst-Allokation).
Das Ganze dauert üblicherweise sub-Sekunde bis ein paar
Sekunden, abhängig davon, wie viele Entities im Shard sind und
wie lang ihre postStop-Arbeit dauert.
Wer einen Handoff anordnen darf
Abschnitt betitelt „Wer einen Handoff anordnen darf“Ein HandOff stoppt jede Entity eines Shards, deshalb prüft die
Region woher er kam, bevor sie ihn ausführt — nicht nur, was
darin steht. Die Regel gilt für jede Direktive, die ausschließlich
der Koordinator aussprechen darf: HandOff, ShardHome,
RememberedEntities, ShardMapUpdate und
RegisterAcknowledgment.
Zwei Bedingungen müssen erfüllt sein:
- Der Frame kam über den eigenen Envelope-Pfad der Region herein, der ihn mit dem Peer stempelt, als der die Verbindung authentifiziert wurde — niemals mit einem Feld aus der Payload, das der Absender selbst schreibt.
- Dieser Peer ist der Node, von dem die Region aktuell glaubt, dass er den Koordinator hostet, also der aktuelle Leader.
Alles andere wird verworfen und auf WARN geloggt. Bis es diese
Regel gab, konnte jeder Peer, der den Cluster-Handshake abschließen
konnte, mit einem kleinen Frame die Entities eines Shards
abräumen — beliebig oft; die Region dispatchte allein auf das
kind der Nachricht und kannte überhaupt keinen Absender.
Ein Verwerfen ist konstruktionsbedingt unkritisch: Eine Region registriert sich beim nächsten Leader- oder Membership-Event neu und fragt jeden gepufferten Shard erneut an — eine während eines Leader-Wechsels verworfene Direktive wird also erneut angefordert statt verloren.
Eine dritte Bedingung betrifft den Shard statt den Absender: Die
Region gibt nur einen Shard ab, den sie aktuell besitzt, und nur
einmal. Ein HandOff, den sie schon ausführt, oder einer, der einen
Shard nennt, der woanders liegt, wird ignoriert — vorher leerte ein
solcher Frame trotzdem die gecachten „Shard X liegt auf Node
N”-Einträge der Region, sodass eine authentische Direktive, die
zweimal ankam oder verspätet, nachdem der Koordinator den Shard
längst zwangsweise neu allokiert hatte, den Verkehr jedes anderen
Nodes einen zusätzlichen Round-Trip zum Koordinator kostete. Besitz
richtet sich nach der Allokation, nicht danach, ob der Shard-Actor
gerade läuft: Ein Shard, der wegen Untätigkeit passiviert wurde,
gehört weiterhin dazu und wird weiterhin abgegeben.
Der Koordinator bleibt dabei nicht hängen. Er schickt HandOff
ausschließlich an die Region, die seine eigene Allokationskarte
nennt — eine Verweigerung heißt also, dass die beiden sich über den
Besitz uneins sind, und genau dafür ist handOffTimeoutMs weiter
unten der Rückfallpfad: Der Shard wird zwangsweise neu allokiert,
statt endlos zu warten.
Konfiguration
Abschnitt betitelt „Konfiguration“const shardingOptions = StartShardingOptions.create() // ... .withRebalanceIntervalMs(2_000) // Strategie-getriebenes Rebalance alle 2s .withHandOffTimeoutMs(10_000); // nach 10s aufgeben + zwangsweise neu allokierensharding.start(shardingOptions);| Knopf | Default | Was |
|---|---|---|
rebalanceIntervalMs | 2000 | Wie oft die Strategie nach Rebalance-Empfehlungen befragt wird. |
handOffTimeoutMs | 10_000 | Wenn die Quellregion nicht innerhalb dieses Fensters HandOffComplete sendet, gibt der Koordinator das Warten auf und allokiert zwangsweise neu. |
Was gepufferte Nachrichten bedeuten
Abschnitt betitelt „Was gepufferte Nachrichten bedeuten“Während Shard X handing-off ist:
- Nachrichten an Entities in
Xkommen bei Node A an. - Die Region von Node A puffert sie — leitet sie nicht an die (jetzt gestoppten) Entity-Actors weiter.
- Sobald der Handoff abgeschlossen ist, fragt Node A den Koordinator,
wohin
Xgegangen ist, und leitet den Puffer an den neuen Besitzer weiter, sobald die Antwort da ist. - Die neue Region spawnt Entities und verarbeitet die Nachrichten in Reihenfolge.
Das bedeutet, dass eine während des Handoffs gesendete Nachricht verzögert ist, nicht verloren. Die Kosten: Latenzspitzen während Rebalance-Fenstern.
Das Nachfragen ist wichtiger, als es aussieht. Ein abgeschlossener Handoff löscht das gecachte Shard-Home der alten Region, und der Koordinator verkündet eine neue Platzierung nur dem neuen Besitzer und Regionen mit offener Anfrage — ohne Nachfrage bliebe die alte Region also auf ihrem Puffer sitzen, bis zufällig eine spätere, unabhängige Nachricht für denselben Shard ein Lookup auslöst. Bei einem Shard, der nach dem Rebalance still wird, kommt die nie.
Wenn Entities State haben
Abschnitt betitelt „Wenn Entities State haben“Für persistente Entities (PersistentActor):
postStopauf der Quellseite schließt etwaiges ausstehendes Persist ab.- Die neue Entity-Instanz auf dem Ziel spielt das Journal beim Start ab.
- Die gepufferte Nachricht kommt bei einer vollständig wiederhergestellten Entity an.
Für nicht-persistente Entities ist der State zwischen Inkarnationen verloren — der Rebalance entspricht einem Restart, bei dem die gepufferte Nachricht als erstes Kommando wirkt.
Wenn State über Rebalance hinweg zählt, nutze Persistenz. Das ist nicht optional — und es hat selbst eine Voraussetzung: das Ziel replayt das Journal, das sein Knoten öffnet. Das Versprechen hält also nur, wenn jeder Knoten dieselbe Datenbank liest. Pro-Knoten-Storage (eine SQLite-Datei, die In-Memory-Defaults) macht aus dem Rebalance stattdessen einen stillen State-Fork — die Entity recovert eine andere Historie, ohne Fehler irgendwo. Das Cluster warnt vor dieser Kombination; siehe Storage-Lokalität & Identität.
Zwangsweise Neuallokation
Abschnitt betitelt „Zwangsweise Neuallokation“Wenn HandOffComplete nicht innerhalb von handOffTimeoutMs
eintrifft:
- Der Koordinator loggt eine Warnung (“handoff timed out”).
- Er allokiert den Shard zwangsweise auf seinen neuen Besitzer.
- Gepufferte Nachrichten leiten weiter, als wäre der Handoff normal abgeschlossen.
- Die alten Entities sind möglicherweise noch am Leben auf Node A — sie empfangen auch lokal Nachrichten, was zu Split-Brain-Entities führt, bis die Waisen sterben.
Das ist selten, aber möglich. Ursachen:
- Ein
postStop, das hängt (langsamer Journal-Write, blockierender externer Call). - Eine Netzwerkpartition während des Handoffs.
Abhilfe:
- Halte
postStopkurz und nicht-blockierend. - Tune
handOffTimeoutMsauf deine Worst-Case-Persist-Zeit + Puffer. - Für das Split-Brain-Risiko erwäge eine Lease auf dem Sharding-Koordinator (siehe Sharding mit Lease).
Rebalance vs. Skalierung
Abschnitt betitelt „Rebalance vs. Skalierung“Node hinzugefügt → Koordinator allokiert ihm einige ShardsNode entfernt → Koordinator findet neue Heimaten für seine ShardsNode ↑ in Last → LeastShardAllocationStrategy drainiert Shards wegNode ↓ in Last → kein Rebalance (nur die belastete Richtung löst aus)Rebalance ist Arbeitsabgabe, nicht Arbeitseinwerbung. Ein
leerlaufender Node bekommt aus Sicht von
LeastShardAllocationStrategy neue Shards, sobald sie allokiert
werden, aber er bekommt keine existierenden Shards in sich
verschoben, nur weil er ruhig ist.
Für einen “immer-balanciert”-Effekt starte den Koordinator
periodisch neu (was Allokation von Null neu laufen lässt) — oder
konfiguriere rebalanceThreshold: 1, um auf jedes Ungleichgewicht
zu reagieren.
Wohin als Nächstes
Abschnitt betitelt „Wohin als Nächstes“- Sharding-Überblick — das größere Bild.
- Allokationsstrategie — was die neuen Heimaten entscheidet.
- Remember Entities — beeinflusst, was nach Handoff neu gespawnt wird.
- Sharding mit Lease — Split-Brain-Schutz für den Koordinator.
