Zum Inhalt springen
Deutsch

DistributedPubSub

DistributedPubSub ist die clusterweite Version des lokalen Event-Streams — Pub/Sub per Topic-Namen, knotenübergreifend.

per Gossip bekannte Subs

auf anderen Nodes

Publisher

Mediator von node-A

lokaler Sub auf node-A

Mediator von node-B

lokaler Sub auf node-B

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.

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.

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

Bei Publish(topic, message):

  1. Der lokale Mediator schlägt das Topic in seiner Map<topic, { local, remoteNodes }> nach.
  2. Lokale Subscriber empfangen direkt — local.values(), jeder bekommt ein tell.
  3. 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.

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.

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.

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));
}
}

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:

OptionHOCON-LeafDefaultBegrenzt
maxSubscribersPerTopiccluster.pub-sub.max-subscribers-per-topic10000Lokale Subscriber auf einem Topic
maxTopicscluster.pub-sub.max-topics10000Verschiedene Topics auf diesem Mediator
maxRemoteNodesPerTopiccluster.pub-sub.max-remote-nodes-per-topic1000Peers, 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 ist

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

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 = off

Vier gute Anwendungsfälle:

  1. Chat / Notifications — mehrere Subscriber (oft auf verschiedenen Nodes), die sich für dasselbe Topic interessieren.
  2. Systemweite Ankündigungen — ein “schema-updated”-Event, auf das jeder Node reagieren soll.
  3. Entkoppelter Fan-out über Nodes — wenn der Publisher nicht wissen sollte, wie viele Subscriber existieren oder wo sie leben.
  4. 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.

Zwei Pub/Sub-Bus-Implementierungen; wähle nach Scope:

BusScopeTopic-Key
Event-StreamEin ActorSystemKlasse (instanceof)
DistributedPubSubClusterweitString-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.

Die DistributedPubSubMediator API-Referenz deckt das vollständige Protokoll ab.