Zum Inhalt springen
Deutsch

Worker-Mesh

JavaScript ist pro Actor-System Single-Threaded. Für Parallelität innerhalb eines einzelnen OS-Prozesses lässt das Worker-Mesh des Frameworks mehrere ActorSystems laufen — eines pro Worker-Thread — die alle am selben Cluster über einen MessageChannel-Transport teilnehmen.

Hauptprozess (einzelner OS-Prozess)

ActorSystem 'main'

Hauptthread

ActorSystem 'w1'

Worker-Thread 1

ActorSystem 'w2'

Worker-Thread 2

ActorSystem 'w3'

Worker-Thread 3

Jeder ist aus Sicht des Clusters ein separater Cluster-Node — Gossip + Mitgliedschaft + Sharding gelten alle. Die Kommunikation zwischen ihnen läuft über In-Process-MessageChannel (keine Serialisierung zu Bytes, kein TCP).

Zwei Hauptszenarien:

  1. CPU-gebundene Parallelität in einem Prozess — actor-ts ist pro System Single-Threaded; Multi-Threading braucht mehrere Systeme. Worker-Mesh verteilt sie.
  2. Isolation innerhalb eines Prozesses — ein Worker, der ausfällt, reißt das Hauptsystem nicht mit. Das gilt für einen nicht gefangenen Throw, eine unbehandelte Rejection und ein Bootstrap, das sich nicht laden lässt: Das Framework abonniert das error-Event des Workers, sodass der Fehler bei restartPolicy landet und nicht auf dem Crash-Pfad des Hosts. Siehe Fehler-Containment dazu, was pro Runtime garantiert ist und was nicht.

Für Multi-Process-Parallelität (separate OS-Prozesse) nutze den regulären Cluster + TCP-Transport. Worker-Mesh ist speziell für den In-Process-Fall.

Broker, Kanäle und den Per-Worker-Handshake von Hand zu verdrahten (siehe Unter der Haube weiter unten) ist genau das, was WorkerCluster automatisiert. Es spawnt einen Worker-Pool aus einem Entrypoint-Modul, führt den Hello/Init/Ready-Handshake aus, registriert jeden Worker an einem gemeinsamen WorkerBroker und startet abgestürzte Worker gemäß einer restartPolicy neu. Das zugrunde liegende Worker-Primitiv wird pro Runtime gewählt — Web Workers auf Bun/Deno, node:worker_threads auf Node — sodass derselbe Code überall läuft.

// main.ts — Hauptthread
import { WorkerCluster, WorkerClusterOptions } from 'actor-ts/worker';
const workerClusterOptions = WorkerClusterOptions.create()
.withWorkers(4)
.withBootstrap(new URL('./worker-node.js', import.meta.url))
.withSystemName('multi-core')
.withBasePort(2552);
const cluster = await WorkerCluster.spawn(workerClusterOptions);
console.log(`Spawned ${cluster.size} workers:`);
for (const address of cluster.addresses) console.log(' -', address.toString());
// ...später, beim Shutdown:
await cluster.terminate();

Jeder Worker führt das Bootstrap-Modul aus. Es ruft WorkerNode.join() auf, um den Handshake abzuschließen und seine Adresse, den System-Namen, den Transport und initData zu erhalten, und tritt dann dem Cluster bei wie jeder andere Knoten:

// worker-node.ts — läuft in jedem Worker
import { ActorSystem } from 'actor-ts';
import { Cluster, ClusterOptions } from 'actor-ts/cluster';
import { WorkerNode } from 'actor-ts/worker';
async function main(): Promise<void> {
const context = await WorkerNode.join<{ seedAddress?: string }>();
const system = ActorSystem.create(context.systemName);
const clusterOptions = ClusterOptions.create()
.withHost(context.self.host)
.withPort(context.self.port)
.withSeeds(context.initData.seedAddress ? [context.initData.seedAddress] : [])
.withTransport(context.transport);
await Cluster.join(system, clusterOptions);
system.spawn(MyActor, 'worker');
context.ready(); // dem Hauptthread signalisieren, dass der Knoten läuft
}
void main();
OptionWas
withBootstrap(url)Worker-Entrypoint-Modul. Pflicht.
withWorkers(n | 'auto')Pool-Größe; 'auto' nutzt die Hardware-Concurrency. Default 'auto'.
withSystemName(name)ActorSystem-Name, den jeder Worker hostet. Default 'worker-cluster'.
withHostname(host)Hostname-Komponente der Adresse jedes Workers. Default 'worker'.
withBasePort(port)Port des ersten Workers; jeder weitere inkrementiert. Default 1.
withInitData(data)Payload, die jedem Worker im join()-Context zugestellt wird (initData). Default null.
withRestartPolicy(p)'always' / 'on-failure' / 'never'. Default 'on-failure'.
withReadyTimeoutMs(ms)Handshake-Timeout pro Worker. Default 10000.
withRestartMinBackoffMs(ms)Wartezeit vor dem ersten Respawn eines abgestürzten Slots. Default 200.
withRestartMaxBackoffMs(ms)Obergrenze für die Respawn-Wartezeit, die pro Versuch verdoppelt wird. Default 10000.
withRestartRandomFactor(f)±-Jitter-Anteil auf jede Respawn-Wartezeit, in [0, 1]. Default 0.2.
withMaxRestarts(n)Restarts pro Slot innerhalb des Fensters, bevor er stillgelegt wird; -1 für unbegrenzt. Default 10.
withRestartWindowMs(ms)Gleitendes Fenster, über das das Restart-Budget zählt; 0 setzt es nie zurück. Default 60000.
withOnWorkerPermanentlyDown(fn)Wird einmal pro Slot aufgerufen, dessen Restart-Budget erschöpft ist. Default: eine console.error-Zeile.
withBackend(backend)Startet über dieses WorkerBackend statt über das erkannte — für eine Runtime, die die Erkennung nicht kennt, oder ein In-Memory-Fake im Test. Default: erkannt.

Das zurückgegebene WorkerCluster bietet size, addresses (die NodeAddress jedes Workers), den gemeinsamen broker und terminate(). Fehlkonfigurierte Optionen (ein falscher basePort, ein nicht-positives workers) werfen OptionsError bei spawn.

spawn() lässt bei einem Fehlschlag keine lebenden Worker zurück: Ein Worker, dessen Handshake abläuft, wird beendet, bevor die Rejection weitergegeben wird — und ebenso jedes Geschwister, das erfolgreich hochgekommen ist, denn die Instanz, die sie besitzt, wird nie zurückgegeben, sodass sie sonst niemand mehr erreichen könnte. terminate() wartet darauf, dass die Threads tatsächlich verschwinden, statt nur abzufeuern — innerhalb der unten beschriebenen Schranke.

Das obige WorkerCluster übernimmt diese Verdrahtung für dich. Greife zum manuellen Weg nur, wenn du volle Kontrolle brauchst — eine eigene Topologie, einen maßgeschneiderten Handshake oder zusätzliches Per-Node-Setup.

// main.ts — Hauptthread
import { Worker } from 'node:worker_threads';
import { ActorSystem } from 'actor-ts';
import { Cluster, ClusterOptions, MessageChannelTransport, NodeAddress } from 'actor-ts/cluster';
const channel = new MessageChannel();
const w1 = new Worker('./worker.js', {
workerData: { mainPort: channel.port2 },
transferList: [channel.port2],
});
const mainAddress = new NodeAddress('main', 'main', 0);
const transport = new MessageChannelTransport(mainAddress, channel.port1);
const system = ActorSystem.create('main');
const clusterOptions = ClusterOptions.create()
.withHost('main')
.withPort(0)
.withSeeds(['main'])
.withTransport(transport);
await Cluster.join(
system,
clusterOptions,
);
// worker.js — läuft im Worker-Thread
import { parentPort, workerData } from 'node:worker_threads';
import { ActorSystem } from 'actor-ts';
import { Cluster, ClusterOptions, MessageChannelTransport, NodeAddress } from 'actor-ts/cluster';
const workerAddress = new NodeAddress('w1', 'w1', 0);
const transport = new MessageChannelTransport(workerAddress, workerData.mainPort);
const system = ActorSystem.create('w1');
const cluster2Options = ClusterOptions.create()
.withHost('w1')
.withPort(0)
.withSeeds(['main'])
.withTransport(transport);
await Cluster.join(
system,
cluster2Options,
);
// Ab hier ist w1 einfach ein weiterer Cluster-Node

Der Transport ist ein Stern, kein voll vermaschtes Mesh. Ein einzelner WorkerBroker im Hauptthread ist die Nabe: Er hält einen MessagePort pro Node und leitet jeden Frame an den Node weiter, den die to-Adresse des Envelopes benennt. Der MessageChannelTransport jedes Nodes hält einen einzelnen Port — sein eigenes Ende eines Channels zum Broker, niemals ein Array:

import { MessageChannelTransport, NodeAddress } from 'actor-ts/cluster';
import { WorkerBroker } from 'actor-ts/worker';
// Hauptthread: ein Broker ist die Nabe für jeden Node.
const broker = new WorkerBroker();
// Ein MessageChannel pro Worker — der Broker behält das eine Ende und
// der Worker erhält das andere (übertragen via workerData).
const workerAddress = new NodeAddress('w1', 'w1', 0);
const channel = new MessageChannel();
broker.register(workerAddress, channel.port1);
// In diesem Worker hält sein Transport das einzelne ferne Ende:
const transport = new MessageChannelTransport(workerAddress, channel.port2);

Also brauchen N Worker N Channels (einen pro Worker) — nicht binomial(N,2). Der Broker leitet jeden Hop weiter, was die Verdrahtung linear hält, aber den Hauptthread zum Relay für den gesamten Node-übergreifenden Verkehr macht.

TCP-Transport: MessageChannelTransport:
- Sockets, Framing - postMessage zwischen Threads
- Serialisierte Bytes - Structured Cloning (kein JSON)
- Netzwerklatenz - Sub-Mikrosekunde
- Cross-Host - Nur derselbe Prozess

Nachrichten zwischen Worker-Systemen gehen durch Structured Clone — schneller, als Text zu rahmen und neu zu parsen. Er trägt dieselben reichen Typen (Map, Set, Date …) wie der getaggte JSON-Tree des TCP-Wire, nur ohne den Text-Round-Trip.

// 4-Worker-Mesh; Sharding verteilt Entities über sie:
const startShardingOptions = StartShardingOptions.create()
.withTypeName('order')
.withEntityActor(...)
.withExtractEntityId((message) => message.id)
.withNumShards(16);
sharding.start(
startShardingOptions,
);

Der Koordinator (auf main) allokiert Shards auf die 4 Worker. CPU-gebundene Entity-Arbeit parallelisiert über Cores.

// Worker, der GPU-gebundene Jobs handhabt:
system.spawn(GpuJobActor, 'gpu-jobs');
// Abstürze in diesem Worker bleiben isoliert von main + anderen Workern

Ein abstürzender Worker reißt das Hauptsystem nicht mit — separate Event-Loops.

Was “ein ausfallender Worker reißt das Hauptsystem nicht mit” tatsächlich abdeckt — und was es kostet:

  • Ein nicht gefangener Throw, eine unbehandelte Rejection oder ein Bootstrap, das sich nicht laden lässt, wird abgefangen. Das Framework abonniert das error-Event des Workers auf jeder Runtime; ohne dieses Abonnement wirft Node den Worker-Fehler erneut im Host und Deno rejected eine interne Promise — beide beenden sich mit 1, bevor restartPolicy überhaupt gefragt wird.
  • Ein Fehler erzeugt genau einen Respawn. Node und Bun senden für einen einzigen Throw error und danach das Exit-Event, Deno nur error — was zuerst eintrifft, treibt den Restart, das andere wird verworfen.
  • Respawns warten und haben ein Budget. Jeder Slot wartet restartMinBackoffMs, verdoppelt bis restartMaxBackoffMs mit restartRandomFactor-Jitter, und bekommt maxRestarts Versuche innerhalb von restartWindowMs. Ein Slot, dessen Budget erschöpft ist, wird stillgelegt: Für ihn startet kein weiterer Worker, und onWorkerPermanentlyDown wird einmal mit Index, Adresse und Restart-Zahl des Slots aufgerufen. Ein Ersatz, der nie ready wird, zählt gegen dasselbe Budget — ein dauerhaft defektes Bootstrap hört also auf, statt zu schleifen.
  • Ein fehlgeschlagener Respawn verkleinert das Mesh um einen Worker, statt den Host zu töten. size sinkt, die Adresse des Slots wird beim Broker abgemeldet, und der Fehler wird gemeldet.

Zwei Runtime-Asymmetrien sind erwähnenswert:

Fehlerhafter Traffic wird verworfen, nicht fatal: Der Broker validiert die Adressen to und from eines Frames, bevor er ihn weiterleitet — ein einzelnes fehlerhaftes postMessage aus einem Worker wirft also nicht mehr im Message-Listener des Hosts. Das payload des Frames validiert der empfangende Transport, nicht der Broker.

Der Broker adressiert jeden Frame auf den Port um, über den er hereinkam. Das from eines Frames ist die Adresse, unter der dieser Port registriert wurde — nicht die Adresse, die der sendende Worker hineingeschrieben hat:

// Worker 1 sendet einen Frame und behauptet, Worker 2 zu sein …
port1.postMessage({ from: worker2.toJSON(), to: worker3.toJSON(), payload });
// … und der Transport von Worker 3 meldet trotzdem Worker 1 als Absender.

Diese Adressen vergibt der Host — WorkerCluster erzeugt eine pro Slot, und bei manueller Verdrahtung übergibst du eine an broker.register —, die Peer-Identität eines empfangenden Nodes stammt also aus dem Channel und nicht aus dem Payload. Diese Identität ist es, die sein Failure-Detector gutschreibt, als Absender eines Envelopes gilt und auf die jede Regel nach dem Muster „darf dieser Peer das sagen?” schlüsselt. Ohne die Korrektur könnte ein Worker also ein totes Geschwister lebendig erscheinen lassen und seine eigenen Nachrichten diesem Geschwister zuschreiben lassen.

Korrigiert wird nur der Slot system@host:port. Die optionale incarnation, die eine Adresse tragen kann, wird unverändert durchgereicht, weil im Cluster noch nichts auf sie schlüsselt.

  • Cluster-Überblick — das Cluster-Modell, an dem das Worker-Mesh teilnimmt.
  • Transports — die Transport-Schnittstelle, die MessageChannelTransport implementiert.
  • Sharding — der Hauptnutzer der Mesh-Parallelität.