DistributedPubSub
DistributedPubSub ist die clusterweite Version des lokalen
Event-Streams — Pub/Sub per Topic-Namen, knotenübergreifend.
Jeder Node hostet einen Mediator an einem bekannten Pfad
(/system/cluster/pubsub/mediator). Subscriber registrieren sich bei ihrem
lokalen Mediator; Mediators verteilen die Topic→Node-Karte per
Gossip. Beim Veröffentlichen geht die Nachricht an den lokalen
Mediator, der an jeden Node ausfächert, der Subscriber für dieses
Topic hat.
Ein minimales Beispiel
Abschnitt betitelt „Ein minimales Beispiel“import { ActorSystem, Cluster, ClusterOptions, Actor } from 'actor-ts';import { DistributedPubSubId, type DistributedPubSubMediator, Publish, Subscribe } from 'actor-ts/cluster/pubsub';
class ChatMessage { constructor(public readonly user: string, public readonly text: string) {}}
class ChatRoom extends Actor<ChatMessage> { override onReceive(message: ChatMessage): void { this.log.info(`[chat] ${message.user}: ${message.text}`); }}
const system = ActorSystem.create('my-app');const clusterOptions = ClusterOptions.create() .withHost(host) .withPort(port) .withSeeds(seeds);const cluster = await Cluster.join(system, clusterOptions);const ps = system.extension(DistributedPubSubId);ps.start(cluster);
// Subscribe (typischerweise im preStart eines Actors):const room = system.spawnAnonymous(ChatRoom);ps.mediator.tell(new Subscribe('chat.room.general', room));
// Publish (von überall — von jedem Node, in oder außerhalb eines Actors):ps.mediator.tell(new Publish('chat.room.general', new ChatMessage('alice', 'hi')));Das Publish erreicht jeden Subscriber auf jedem Node — die
Alice-Nachricht kommt bei room an, unabhängig davon, welcher
Node den Publisher hostet.
Die vier Operationen
Abschnitt betitelt „Die vier Operationen“| Nachricht | Was |
|---|---|
Subscribe(topic, ref, replyTo?) | Registriert ref als Subscriber auf topic. Antwortet replyTo (oder dem Sender) mit SubscribeAcknowledgment — oder mit SubscribeRejected, wenn eine Grenze voll ist. |
Unsubscribe(topic, ref) | Entfernt ref aus den Subscribern von topic. |
UnsubscribeAll(ref) | Entfernt ref aus jedem Topic. |
Publish(topic, message, delivery?) | Schickt message an jeden Subscriber von topic — oder, mit delivery = 'one-subscriber', an genau einen. |
Schicke diese an ps.mediator (ein ActorRef). Nutze ask, wenn
du das Ack brauchst:
import {} from 'actor-ts';
await ps.mediator.ask(new Subscribe('chat.room.general', room));Subscribe nimmt ein optionales drittes Argument, replyTo, das
benennt, wohin Bestätigung oder Ablehnung gehen. Ohne es folgt die
Antwort context.sender — und der ist bei der obigen Form
mediator.tell(…), aufgerufen außerhalb eines Actors, leer.
Benenne ein replyTo, wann immer du die Antwort sehen willst:
ps.mediator.tell(new Subscribe('chat.room.general', room, room));Topic-Namen sind beliebige Strings. Das Framework legt keine
Struktur fest — chat.room.general, user-42.events,
metrics-tier-1 funktionieren alle.
Zur Organisation funktioniert eine punkt-segmentierte Konvention
gut (<domain>.<scope>.<resource>), aber das Framework
interpretiert die Segmente nicht — es ist nur String-Matching.
Wie der Fan-out funktioniert
Abschnitt betitelt „Wie der Fan-out funktioniert“Bei Publish(topic, message):
- Der lokale Mediator schlägt das Topic in seiner
Map<topic, { local, remoteNodes }>nach. - Lokale Subscriber empfangen direkt —
local.values(), jeder bekommt eintell. - Remote-Nodes mit Subscribern bekommen einen Envelope pro Node (nicht pro Subscriber) — der Mediator auf dem Zielnode fächert an seine Lokalen aus.
Das ergibt eine höchstens-ein-Remote-Hop-Zustellung: ein Publish kettelt nie über mehrere Nodes, um einen Subscriber zu erreichen.
Anycast — ein Subscriber statt aller
Abschnitt betitelt „Anycast — ein Subscriber statt aller“Ein drittes Argument schaltet Publish von Broadcast auf
Anycast um: genau ein Subscriber im Cluster bekommt die
Nachricht.
// Jeder Worker tritt demselben Topic bei…ps.mediator.tell(new Subscribe('jobs', worker));
// …und jede Aufgabe wird einmal erledigt, von einem von ihnen.ps.mediator.tell(new Publish('jobs', new RenderThumbnail(id), 'one-subscriber'));Das ist die Work-Queue-Form: N Worker über den Cluster verteilt,
jede Aufgabe genau einmal erledigt, und kein Worker muss wissen,
wie viele andere es gibt. Der Default bleibt 'all-subscribers',
ein Publish ohne drittes Argument ist also unverändert.
Wie der eine gewählt wird. Der Mediator baut eine Kandidatenliste — jeder lokale Subscriber zählt einmal, jeder Remote-Node, der das Topic beansprucht, zählt einmal — und läuft sie reihum durch, ein Schritt pro Publish. Zwei Folgen sind es wert, gekannt zu werden:
- Es ist eine Rotation, keine Ziehung. Zehn Aufgaben auf drei Worker landen 4/3/3, nicht “vermutlich ungefähr gleichmäßig”. Ein Neustart des Publishers startet die Rotation neu; sie ist Per-Topic-Zustand des Mediators, keine clusterweite Sequenz.
- Die Kandidatenreihenfolge ist stabil. Zuerst die lokalen Subscriber in Registrierungsreihenfolge, dann die Remote-Anspruchsteller nach Adresse sortiert. Eintreffender Gossip würfelt den Durchlauf nicht neu, und jeder Mediator ordnet die Remote-Hälfte gleich.
- Remote-Nodes zählen als je ein Kandidat, nicht als einer pro Subscriber. Ein Node mit zehn Subscribern bekommt denselben Anteil wie einer mit einem, denn der Gossip-Frame trägt Topic-Namen und keine Subscriber-Zahlen (genau das hält ihn klein). Innerhalb des gewählten Nodes wählt dessen eigene Rotation den Subscriber. Halte die Worker-Zahlen über die Nodes hinweg ungefähr gleich, wenn die Balance zählt.
- Ein Node, der publiziert und Worker hostet, rotiert über zwei Zeiger. Die Anycasts, die er selbst auslöst, laufen über die volle Kandidatenliste; die, die ihn über die Leitung erreichen, nur über seine lokalen Subscriber. Die beiden Listen sind unterschiedlich lang und werden deshalb getrennt gezählt — ein gemeinsamer Zeiger würde den auslösenden Durchlauf unter der Zahl der lokalen Subscriber festnageln und jeden Remote-Kandidaten aushungern.
Ein Anycast, das keinen Kandidaten findet — kein lokaler Subscriber und kein Remote-Anspruch —, geht wie jedes andere ungeroutete Publish an die Dead Letters. Genauso eines, das einen Hop zurückgelegt hat und dessen Subscriber auf der Gegenseite schon weg sind: es wird nicht umgeroutet, denn ein zweiter Hop würde die Höchstens-ein-Hop-Garantie gegen ein Rennen mit der Gossip-Runde tauschen, die den Sender ohnehin gleich korrigiert.
Auch ein Frame, für den der Mediator überhaupt keinen Handler hat — etwa ein älterer Node, den eine neuere Wire-Art erreicht —, landet in den Dead Letters, mit einer Warnung. Rolling Upgrades sind der Fall, der das braucht: ein still verworfener Frame sieht genau aus wie ein Cluster, der nichts zu tun hat.
Topic→Node-Gossip
Abschnitt betitelt „Topic→Node-Gossip“Der Mediator hält seine Map<topic, SubscriberSet> lokal, aber
gossipt Deltas zu Peers:
- “Node X hat jetzt Subscriber für Topic Y.”
- “Node X hat keine Subscriber mehr für Topic Y.”
Standard-Gossip-Intervall ist gossipIntervalMs des Clusters
(1 Sekunde). Überschreibe pro Mediator:
const distributedPubSubOptions = DistributedPubSubOptions.create().withGossipIntervalMs(500);system.extension(DistributedPubSubId).start(cluster, distributedPubSubOptions);Niedrigere Intervalle → schnellere Konvergenz nach Subscribe / Unsubscribe, mehr Geplapper. 500 ms ist für Chat-artige Anwendungen vernünftig.
Wenn Subscriber stoppen
Abschnitt betitelt „Wenn Subscriber stoppen“Der Mediator watcht jeden lokalen Subscriber. Ein Subscriber,
der stoppt, ohne Unsubscribe zu senden — Crash, vergessenes
Cleanup, ein pro Request gespawnter Actor —, wird entfernt, sobald
das Terminated ankommt; mit seinem letzten Subscriber fällt auch
das Topic weg.
Explizit zu unsubscriben bleibt der schnellere Weg und ist die einzige Möglichkeit, ein Topic zu verlassen, ohne zu stoppen:
class Subscriber extends Actor<...> { override preStart(): void { this.system.extension(...).mediator.tell(new Subscribe('topic', this.self)); } override postStop(): void { this.system.extension(...).mediator.tell(new UnsubscribeAll(this.self)); }}Begrenzte Registries
Abschnitt betitelt „Begrenzte Registries“Der Mediator hält drei Dinge, die ein Subscriber — oder der Gossip eines Peers — sonst unbegrenzt wachsen lassen könnte, und der Publish-Fan-out läuft über alle drei. Jedes hat eine Obergrenze:
| Option | HOCON-Leaf | Default | Begrenzt |
|---|---|---|---|
maxSubscribersPerTopic | cluster.pub-sub.max-subscribers-per-topic | 10000 | Lokale Subscriber auf einem Topic |
maxTopics | cluster.pub-sub.max-topics | 10000 | Verschiedene Topics auf diesem Mediator |
maxRemoteNodesPerTopic | cluster.pub-sub.max-remote-nodes-per-topic | 1000 | Peers, die Subscriber für ein Topic beanspruchen |
const distributedPubSubOptions = DistributedPubSubOptions.create() .withMaxSubscribersPerTopic(500) .withMaxTopics(1_000);system.extension(DistributedPubSubId).start(cluster, distributedPubSubOptions);Ein Subscribe über maxSubscribersPerTopic oder maxTopics wird
abgelehnt, nicht verworfen — die Antwort ist ein
SubscribeRejected mit der Grenze, die es abgelehnt hat:
import { SubscribeRejected } from 'actor-ts';
// message.reason: 'maxSubscribersPerTopic' | 'maxTopics'// message.limit: der Wert, auf den diese Grenze gesetzt istmaxTopics und maxRemoteNodesPerTopic gelten auch für
Gossip, und das ist die Hälfte, die man kennen sollte: ein Peer,
der 100 000 Topics ankündigt, für die er angeblich Subscriber hat,
hat früher auf jedem empfangenden Node 100 000 Einträge angelegt —
ganz ohne ein lokales Subscribe. Ansprüche über einer Grenze
werden verworfen und geloggt; die Verbindung bleibt bestehen, denn
ein lauter Peer soll keine gesunde Verbindung kosten.
Publishes, die niemanden erreichen
Abschnitt betitelt „Publishes, die niemanden erreichen“Ein Publish auf ein Topic ohne lokale Subscriber und ohne
Remote-Ansprüche geht an system.deadLetters — genauso wie eines,
das einen Hop zurückgelegt hat und am anderen Ende keinen
Subscriber fand.
system.eventStream.subscribe(monitor, DeadLetter);Das ist standardmäßig an. Ein vertippter Topic-Name und ein Topic, dessen Subscriber noch nicht durchgegossipt sind, sehen von der Publisher-Seite identisch aus — und genau das unterscheidet die beiden. Für ein Deployment, in dem unzustellbare Publishes erwartet werden und die Menge nur Rauschen wäre, schaltest du es ab:
actor-ts.cluster.pub-sub.send-to-dead-letters-when-no-subscribers = offWann DistributedPubSub einsetzen
Abschnitt betitelt „Wann DistributedPubSub einsetzen“Vier gute Anwendungsfälle:
- Chat / Notifications — mehrere Subscriber (oft auf verschiedenen Nodes), die sich für dasselbe Topic interessieren.
- Systemweite Ankündigungen — ein “schema-updated”-Event, auf das jeder Node reagieren soll.
- Entkoppelter Fan-out über Nodes — wenn der Publisher nicht wissen sollte, wie viele Subscriber existieren oder wo sie leben.
- Work-Queues über eine dynamische Worker-Menge — mit
'one-subscriber'-Delivery treten Worker zur Laufzeit einem Topic bei und verlassen es, und jede Aufgabe wird einmal erledigt.
Wann NICHT
Abschnitt betitelt „Wann NICHT“DistributedPubSub vs. Event-Stream
Abschnitt betitelt „DistributedPubSub vs. Event-Stream“Zwei Pub/Sub-Bus-Implementierungen; wähle nach Scope:
| Bus | Scope | Topic-Key |
|---|---|---|
| Event-Stream | Ein ActorSystem | Klasse (instanceof) |
DistributedPubSub | Clusterweit | String-Topic |
Nutze den Event-Stream für In-System-Dispatch; nutze DistributedPubSub, wenn Topics Node-Grenzen überschreiten. Beide können koexistieren — viele Apps nutzen beide für verschiedene Belange.
Wohin als Nächstes
Abschnitt betitelt „Wohin als Nächstes“- Cluster-Überblick — die Mitgliedschaft darunter.
- Event-Stream — das Single-System-Pub/Sub zum Vergleich.
- Refs über Nodes hinweg — wie die cross-node Zustellungen des Mediators serialisieren.
Die DistributedPubSubMediator
API-Referenz deckt das vollständige Protokoll ab.
