Zum Inhalt springen
Deutsch

Sharding — Überblick

Cluster-Sharding ist die Antwort des Frameworks auf “ich habe ein paar Millionen Entities, jede braucht ihren eigenen Actor, verteilt über N Nodes.” Beispiele: Per-User-Sessions, Per-IoT-Device-Controller, Per-Order-Koordinatoren, Per-Game-Room-Actors.

Der Nutzer gibt dem Framework eine Möglichkeit, eine Entity-ID aus jeder Nachricht zu extrahieren; das Framework hasht die ID auf einen Shard; der Koordinator entscheidet, welcher Node diesen Shard hostet; die lokale Region des Nodes routet an diesen Shard, der den Entity-Actor bei Bedarf spawnt. Wenn Nodes kommen und gehen, rebalanciert der Koordinator die Shards über die neue Topologie.

Cluster aus 3 Nodes

node-1 Region

node-2 Region

node-3 Region

Shards 1-33

Shards 34-66

Shards 67-100

Entities für Keys,

die in diese Shards hashen

Entities für Keys,

die in diese Shards hashen

Entities für Keys,

die in diese Shards hashen

Jede Entity ist ein Actor, auf einem Node, zu einer Zeit — genauso wie ein Singleton, aber skaliert auf N Entities. Failover passiert automatisch: wenn ein Node geht, wandern seine Shards woandershin und die Entities werden dort neu gespawnt.

import { match } from 'ts-pattern';
import { Actor, Cluster, ClusterBootstrapOptions, StartShardingOptions } from 'actor-ts';
type AddCommand = { entityId: string; kind: 'add'; sku: string };
type ViewCommand = { entityId: string; kind: 'view'; replyTo: ActorRef<Cart> };
type CartCommand = AddCommand | ViewCommand;
class CartActor extends Actor<CartCommand> {
private items: string[] = [];
override onReceive(command: CartCommand): void {
match(command)
.with({ kind: 'add' }, (c) => this.onAdd(c))
.with({ kind: 'view' }, (c) => this.onView(c))
.exhaustive();
}
private onAdd(command: AddCommand): void { this.items.push(command.sku); }
private onView(command: ViewCommand): void { command.replyTo.tell({ items: this.items }); }
}
// Setup — einzeiliger Cluster- + Sharding-Einstieg:
const { system, cluster } = await Cluster.bootstrap(ClusterBootstrapOptions.create('my-app'));
const startShardingOptions = StartShardingOptions.create<CartCommand>().withExtractEntityId((message) => message.entityId);
const cartRegion = cluster.sharding.start('cart', CartActor,
startShardingOptions);
// Verwendung — `tell` an die Region, mit Entity-ID in der Nachricht:
cartRegion.tell({ entityId: 'user-42', kind: 'add', sku: 'book-1' });
cartRegion.tell({ entityId: 'user-42', kind: 'view', replyTo: ... });
// ^^^^^^^^^^
// Gleiche ID → gleicher Entity-Actor, jedes Mal, unabhängig vom Node.

cluster.sharding.start() akzeptiert mehrere Aufruf-Formen — nimm die, die zur Entity passt. Die kürzeste lässt die Entity ihre eigene Identität deklarieren, wodurch weder Type-Name noch Entity-ID-Extractor am Aufrufort wiederholt werden:

class CartActor extends Actor<CartCommand> {
static readonly shard = ShardKey.of<CartCommand>('cart', (command) => command.entityId);
// ...
}
const cartRegion = cluster.sharding.start(CartActor);
const cart = cluster.sharding.entityRefFor(CartActor, 'user-42');

Die Identität eines ShardKey ist allein sein typeName; der Extractor reitet mit, damit die deklarierende Klasse die einzige Wahrheitsquelle ist — ein Node, der Entities nur nachschlägt, kann denselben Typ mit ShardKey.of<CartCommand>('cart') ohne Extractor benennen. Ein extractEntityId in den Options überschreibt den auf dem Key.

Die übrigen Formen bleiben unverändert:

// 1. Klassen-Kurzform:
const startShardingOptions = StartShardingOptions.create<CartCommand>().withExtractEntityId((m) => m.entityId);
cluster.sharding.start('cart', CartActor, startShardingOptions);
// 2. Factory-Kurzform — wenn die Entity Konstruktor-Argumente braucht:
const startSharding2Options = StartShardingOptions.create<CartCommand>().withExtractEntityId((m) => m.entityId);
cluster.sharding.start('cart', () => new CartActor(deps),
startSharding2Options);
// 3. Vollform — alle Settings über den Builder:
const startSharding3Options = StartShardingOptions.create<CartCommand>()
.withTypeName('cart')
.withEntityActor(CartActor)
.withExtractEntityId((m) => m.entityId)
.withNumShards(16)
.withRole('cart-host');
cluster.sharding.start(
startSharding3Options,
);

cluster.sharding ist eine memoisierte Fassade — wiederholte Zugriffe liefern dieselbe ClusterSharding-Instanz. Die explizite Form ClusterSharding.get(system, cluster) funktioniert weiter und liefert dasselbe Objekt; greif dazu nur, wenn Du die Klasse außerhalb eines Cluster-Handles brauchst.

Aus Sicht des Aufrufers ist die cartRegion ein einzelner ActorRef. Hinter den Kulissen:

  1. Die Region berechnet aus entityId einen Shard (Default: Hash auf einen von 64 Shards).
  2. Sie fragt den Koordinator “wem gehört dieser Shard?”
  3. Wenn der Besitzer dieser Node ist, spawnt sie die Entity (falls noch nicht vorhanden) und leitet die Nachricht weiter.
  4. Wenn der Besitzer ein anderer Node ist, leitet sie über den Cluster-Transport an dessen Region, die dasselbe macht.
ActorRolle
RegionEine pro Node. Routet Nachrichten an den Besitzer des richtigen Shards und entscheidet, wann lokale Entities — und leer gewordene Shards — passiviert werden.
ShardEiner pro hier gehostetem Shard, Kind der Region. Besitzt die Entity-Actors: spawnt sie, stoppt sie, puffert für eine, die gerade geht. Verschwindet selbst, sobald er sein Fenster lang leer stand, und kommt mit der nächsten Nachricht zurück.
KoordinatorEiner pro Cluster (Singleton, auf dem Leader). Entscheidet, welcher Node welchen Shard besitzt. Behandelt Rebalancing bei Mitgliedschaftsänderungen.
EntityEine pro entityId. Kind ihres Shards, auf dem Node, dem aktuell der Shard gehört, in den die ID hasht.

Damit haben Entities einen echten Pfad — /system/cluster/sharding/region-counter/shard-7/entity-user-42 — und genau das macht einen Shard überhaupt adressierbar. Siehe Introspektion zum Auflisten von Shards und zum Holen von Refs darauf.

Die Region ist das Thema dieser Seite; die Entscheidungen des Koordinators haben eigene Seiten — siehe die Allokationsstrategie und Rebalance für die Mechanik.

Eine Entity braucht fast immer die ID, über die sie geroutet wurde — für den Namen ihres Journal-Streams, ihres Pub/Sub-Topics oder einer Zeile in einer fremden Datenbank. Sie liest sie an sich selbst ab:

class CartEntity extends Actor<CartCommand> {
override preStart(): void {
this.log.info(`cart ${this.entityId} woke up in shard ${this.entity.shardId}`);
}
}
this.entityIdDer Wert, den extractEntityId zurückgegeben hat — unverändert.
this.entityDieselbe ID plus typeName und shardId.
this.context.entityDie Option-Form — None, wenn dieser Actor keine Entity ist. Die beiden Getter darüber werfen stattdessen.

Der häufigste Fall ist eine persistenceId pro Entity — genau das gibt einem geshardeten PersistentActor einen Event-Stream pro Entity statt einen pro Typ:

class CartEntity extends PersistentActor<CartCommand, CartEvent, CartState> {
override get persistenceId(): string { return `cart-${this.entityId}`; }
}

Die Identität sitzt auf der Entity und sonst nirgends — die eigenen Kinder einer Entity bekommen None, this.context.entity.nonEmpty beantwortet also „bin ich die Entity?“. Ein Kind, das die ID braucht, bekommt sie explizit weitergereicht.

Um eine Entity isoliert zu testen, gibst du ihr die Identität direkt, statt einen Cluster darum herum hochzufahren:

import { ActorOptions } from 'actor-ts';
const cartOptions = ActorOptions.create()
.withEntity({ entityId: 'user-42', typeName: 'cart', shardId: 3 });
const cart = system.spawn(CartEntity, 'cart-under-test', cartOptions);

.../shard-7/entity-user-42 sieht aus, als trüge er die ID, und für einfache IDs funktioniert das Abschneiden des entity--Präfixes auch. Es ist aber nicht die ID: Actor-Namen haben ein eingeschränktes Alphabet, deshalb maskiert der Shard beim Benennen des Kindes alles, was ein Name nicht tragen kann — eine Code-Unit außerhalb von [A-Za-z0-9_-.@:+] wird zu ~ plus vier Hex-Ziffern, user/42 heißt also entity-user~002F42.

Die Maskierung ist injektiv: zwei verschiedene IDs landen nie auf einem Namen. Das ist eine Korrektheitseigenschaft, keine Feinheit — ein geteilter Name bedeutet, dass die zweite Entity unter einem schon vergebenen Namen startet, und der resultierende Fehler reißt den Shard mitsamt allen anderen Entities darin mit.

Was sie nicht ist: eine zweite Schreibweise der ID. Nichts dekodiert ein Pfadsegment, bewusst, und die Maskierung ist kein Teil der API. Der Pfad ist ein Etikett; entityId ist der Wert.

ShardingOptionsType<TMessage> — die Felder, die du am häufigsten anfasst:

type ShardingOptionsType<TMessage> = {
typeName: string;
entityActor: ActorClassOrFactory<TMessage>;
entityOptions?: ActorOptions<TMessage>;
extractEntityId: (message: TMessage) => string;
extractEntityMessage?: (message: TMessage) => unknown;
numShards?: number; // Default 64
role?: string; // auf Nodes mit dieser Rolle beschränken
proxy?: boolean; // nur Routing, keine lokalen Entities
rememberEntities?: boolean; // Entities bei Failover neu spawnen (#)
passivationIdleMs?: number; // idle Entities auto-stoppen, Default 5 min
shardPassivationIdleMs?: number; // leere Shards auto-stoppen, folgt obigem
maxEntities?: number; // LRU-Cap pro Node
};
FeldWas es steuert
typeNameEin String, der diesen sharded Typ identifiziert. Verschiedene Typen können im selben Cluster koexistieren (cart, session, order).
extractEntityId(message)Zieh die Entity-ID aus einer Nachricht. Das ist der Key, der zu einem Shard gehasht wird.
extractEntityMessage(message)(Optional) Wenn der Nachrichten-Envelope Routing-Info plus eine Payload enthält, entfernt das den Envelope, bevor die Entity ihn sieht. Default ist die Nachricht unverändert.
numShardsIn wie viele Shards der Entity-Raum aufgeteilt wird. 64 reichen für die meisten Cluster; setze auf 1000 für sehr große Cluster (>50 Nodes).
roleNur Mitglieder mit dieser Rolle hosten Shards dieses Typs. Nützlich, um compute-schwere Entities auf dedizierten Nodes zu platzieren.
proxyDieser Node leitet Nachrichten an die Region weiter, hostet aber nie lokale Entities. Wird für Client-Nodes in einem asymmetrischen Cluster verwendet.
rememberEntitiesPersistiere die Menge der aktiven Entity-IDs. Nach einem Koordinator-Failover (oder vollem Cluster-Neustart) werden diese IDs eager gespawnt, sodass Nachrichten nicht die gesamte Flotte neu erstellen müssen.
passivationIdleMsStoppt eine Entity nach so viel Leerlaufzeit. Befreit Speicher; die nächste Nachricht für dieselbe ID erstellt die Entity neu. Default sind 5 Minuten; 0 schaltet es ab.
shardPassivationIdleMsStoppt einen Shard, nachdem er so lange leer stand. Ungesetzt folgt er passivationIdleMs; 0 hält leere Shards resident, während Entities weiterhin passivieren.
maxEntitiesPer-Node-Cap. Wird er überschritten, wird die LRU-Entity passiviert.

Die Defaults sind sinnvoll für kleine Cluster. Für Produktion willst du meistens rememberEntities: true und ein passivationIdleMs passend zu deinem Verkehrsmuster.

Fünf davon sind auch HOCON-Keys — number-of-shards, remember-entities, passivation-idle, shard-passivation-idle und max-entities — ebenso die Koordinator-Keys rebalance-interval und hand-off-timeout. Sie setzen die node-weite Grundeinstellung für jeden sharded Typ und liegen unter dem, was start(...) explizit übergibt — ein Per-Typ-Wert im Builder gewinnt also weiterhin. Siehe Konfiguration.

actor-ts.sharding {
number-of-shards = 128
passivation-idle = 2 minutes
max-entities = 50000
}

Passivierung stoppt eine idle Entity, um Speicher freizugeben; die nächste Nachricht für dieselbe ID erstellt sie neu. Wie auch immer sie ausgelöst wird — die Region puffert Nachrichten, die eintreffen, während die Entity stoppt, und spielt sie an die frische Inkarnation zurück, sodass nichts verloren geht.

Das Standard-Idle-Fenster beträgt 5 Minuten. Setze passivationIdleMs — im Code oder über actor-ts.sharding.passivation-idle —, um es zu ändern, oder 0, um es abzuschalten und jede Entity für die Lebensdauer ihres Nodes resident zu halten.

Automatisch — wenn passivationIdleMs ohne Verkehr abläuft oder maxEntities überschritten wird (die LRU-Entity wird verdrängt), wird die Entity gestoppt. Es wird nichts an die Entity geschickt, und es gibt keine Quittung. Die Entscheidung trifft die Region, weil nur sie den Verkehr über alle Shards dieses Nodes sieht — deshalb ist maxEntities eine Obergrenze pro Node, nicht pro Shard — und der Shard, dem die Entity gehört, führt sie aus.

Manuell (graceful) — eine Entity bittet darum, passiviert zu werden, indem sie ein Passivate an ihren Parent schickt, also an ihren Shard. Passivate trägt eine Stop-Nachricht, die an die Entity zurückgeleitet wird, sodass die Entity genau entscheidet, wann und wie sie herunterfährt (laufende Arbeit abschließen, State flushen, dann terminieren). PoisonPill als Stop-Nachricht terminiert die Entity, sobald sie zurückkommt:

import { match } from 'ts-pattern';
import { Passivate, PoisonPill, Actor } from 'actor-ts';
type AddCommand = { entityId: string; kind: 'add'; sku: string };
type CheckoutCommand = { entityId: string; kind: 'checkout' };
type CartCommand = AddCommand | CheckoutCommand;
class CartActor extends Actor<CartCommand> {
private items: string[] = [];
override onReceive(command: CartCommand): void {
match(command)
.with({ kind: 'add' }, (c) => this.onAdd(c))
.with({ kind: 'checkout' }, () => this.onCheckout())
.exhaustive();
}
private onAdd(command: AddCommand): void {
this.items.push(command.sku);
}
// Mit dem Warenkorb fertig — die Region bitten, uns zu passivieren.
private onCheckout(): void {
this.context.parent.forEach((parent) =>
parent.tell(new Passivate(PoisonPill.instance, this.self), this.self));
}
}

Alle Entities eines Shards zu passivieren ließ früher den Shard selbst zurück — einen lebenden Actor mit einer leeren Map. Außer einem Rebalance stoppte ihn nichts, und da Entity-IDs über den Hash-Raum streuen, wurde auf einem lang laufenden Node irgendwann jeder Shard einmal angefasst und blieb dann. Bei numShards = 64 sind das 64 idle Actors ohne Inhalt.

Deshalb wird auch ein Shard gestoppt, der shardPassivationIdleMs lang leer stand. Die Region behält die Ownership — nur der Actor geht —, der Shard bleibt also routbar, und die nächste Nachricht für ihn erstellt Shard und Entity gemeinsam neu, ohne Umweg über den Koordinator. Ungesetzt folgt das Fenster passivationIdleMs, was meistens genau richtig ist: ein Shard ist leer, gerade weil seine Entities idle wurden.

const shardingOptions = StartShardingOptions.create<CartCommand>()
.withTypeName('cart')
.withEntityActor(CartActor)
.withExtractEntityId((command) => command.entityId)
.withPassivationIdleMs(120_000)
// Shards sind billig zu halten und kosten einen Respawn, wenn sie gehen —
// auf einem belebten Node ist es ein guter Handel, sie ihre Entities eine
// Weile überleben zu lassen.
.withShardPassivationIdleMs(600_000);

Setze shardPassivationIdleMs: 0, um leere Shards resident zu halten, während Entities weiterhin passivieren. Region und Koordinator sind in beiden Fällen unberührt — sie laufen für die Lebensdauer des Nodes.

Gestoppt wird immer nur ein leerer Shard, es steht also kein Entity-State auf dem Spiel, und Nachrichten, die während des Stoppens eintreffen, werden genauso gepuffert und nachgespielt wie bei einer Entity. Auch eine Shard-Ref bleibt über die Lücke hinweg gültig — siehe Introspection, wo ShardInfo zusätzlich ein resident-Flag meldet, das einen passivierten Shard von einem bloß leeren unterscheidet.

Wenn sich die Cluster-Topologie ändert (ein Node tritt bei oder geht), führt der Koordinator einen Rebalance-Lauf durch:

  1. Berechne die neue Shard-zu-Node-Zuordnung aus der aktiven Allokationsstrategie (Default: Hash modulo Regionen).
  2. Sage für jeden bewegten Shard der alten Region, sie soll den Shard übergeben.
  3. Die alte Region stoppt ihre Entities (die ihren State persistieren können), sagt dem Koordinator “Handoff abgeschlossen” und hört auf, für diesen Shard zu routen.
  4. Die neue Region spawnt Entities für diesen Shard bei Bedarf, wenn Nachrichten ankommen.

Der Handoff ist nicht sofort fertig — gepufferte Nachrichten warten auf das “Handoff abgeschlossen”-Signal, bevor sie an den neuen Besitzer weitergeleitet werden. Das verhindert, dass Nachrichten an einer halb-verschobenen Entity vorbeisausen.

Siehe Rebalance für das vollständige Protokoll.

Sharded Entities unterliegen denselben Neustart-Semantiken wie jeder andere Actor — wenn eine Entity zu einem neuen Node wandert, startet die neue Instanz mit weißer Weste.

Für State, der überleben soll:

  • PersistentActor — die Entity persistiert Events in einem Journal; auf einem frischen Node spielt sie das Journal beim Start ab. Siehe PersistentActor.
  • DurableStateActor — einfacher: persistiere den aktuellen State-Snapshot; restore beim Neustart. Siehe DurableState.
  • DistributedData — für State, der von jedem Node lesbar sein soll (nicht nur vom aktuellen Host der Entity), nutze ein CRDT im DD-Replicator stattdessen.

Ohne eine dieser Lösungen gibt dir Sharding Platzierung und Routing, aber keine Durability.

Drei gute Anwendungen:

  1. Per-User-/Per-Tenant-State, der zu viel für einen Node ist, aber nicht von überall gleichzeitig lesbar sein muss.
  2. Per-Entity-Workflows — Sagas, Order-Processing, lange laufende Koordinatoren — die von Per-Key-Serialisierung profitieren.
  3. Hotspots, die Keys folgen — eine Streaming-Pipeline, in der die Events jedes Users in Reihenfolge auf einem einzelnen Actor verarbeitet werden sollen.

Die API-Referenzen ClusterSharding und ShardRegion decken die vollständige Oberfläche ab, inklusive des Per-Node-Region-Actors.