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, Actor } from 'actor-ts';
import { Cluster, ClusterOptions } from 'actor-ts/cluster';
import { DistributedPubSubId, type DistributedPubSubMediator, Publish, Subscribe } from 'actor-ts/cluster';
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:

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 Topic-Fan-out trägt keinen Sender: tell ist einseitig, der Publisher ist meist nicht die interessante Partei, und jedes onReceive eines Subscribers ist gegen den nackten Body geschrieben. Ein Subscriber, der auf den Publisher autorisieren muss, braucht mehr und fordert es beim Subscribe an:

mediator.tell(new Subscribe('ledger', auditor, null, /* deliverWithOrigin= */ true));

Jede Nachricht für diesen Subscriber trifft dann als PubSubEnvelope statt als nackter Body ein:

override onReceive(message: PubSubEnvelope<LedgerEntry>): void {
const { topic, message: entry, origin } = message;
if (origin === null) return; // kein authentifizierter Sender
if (!this.mayWrite(entry.author, origin)) return;
this.apply(entry);
}

origin ist nie ein Wert aus dem Payload. Für eine Nachricht, die über die Leitung kam, ist es der Peer der Verbindung — die Identität aus dem Cluster-Handshake —, für eine auf diesem Node publizierte ist es cluster.selfAddress. null heißt, dass der Mediator keine authentifizierte Identität anhängen konnte; ein Subscriber, der auf origin autorisiert, muss das als nicht authentifiziert lesen und nicht als lokal.

Der Vergleich ist nur wegen der Höchstens-ein-Hop-Regel oben aussagekräftig: der Node, den ein Subscriber sieht, ist der Node, der publiziert hat, und kein Relay, das für jemand anderen weitergeleitet hat.

Das Flag gehört zum Subscriber, nicht zum Abonnement — ein Actor hat ein onReceive, eine Mischung von Formen über seine Topics wäre also ein Unterscheidungsproblem, das er nicht lösen kann. Ein späteres Subscribe desselben Actors legt die Form für alle seine Topics erneut fest. ReplicatedEventSourcedActor ist der Nutzer im Repo: ein Event-Envelope nennt seinen eigenen Autor, und ohne den publizierenden Node gibt es nichts, wogegen sich dieser Name halten ließe.

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/cluster';
// 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.