NATS
NatsActor integriert mit NATS — dem leichtgewichtigen Pub/Sub-
Broker. Core-NATS ist Fire-and-Forget ohne Durability; für
durable Streams siehe die JetStream-Variante.
import { ActorSystem, NatsActor, NatsOptions } from 'actor-ts';
const natsOptions = NatsOptions.create() .withServers(['nats://nats-1:4222', 'nats://nats-2:4222']) .withName('my-app');const nats = system.spawn( () => new NatsActor( natsOptions, ), 'nats',);
// Subscribieren (Wildcards unterstützt):nats.tell({ kind: 'subscribe', subject: 'events.>', target: eventHandler,});
// Publishen:nats.tell({ kind: 'publish', publish: { subject: 'events.user.signup', payload: JSON.stringify(event), },});Settings
Abschnitt betitelt „Settings“interface NatsOptionsType extends BrokerCommonOptionsType { servers: string[] | string; name?: string; // Client-Identifier user?: string; password?: string; token?: string; subscriptions?: { subject: string; target: ActorRef<NatsMessage> }[]; // beim Connect eingerichtet}Subscriptions und Reconnects
Abschnitt betitelt „Subscriptions und Reconnects“subscriptions und das subscribe-Kommando zur Laufzeit speisen dieselbe
gewünschte Menge, die von der
BrokerActor-Basisklasse
gehalten wird und nicht von der Verbindung. Praktische Folgen:
- Jede Subscription — konfiguriert oder zur Laufzeit hinzugefügt — wird nach einem Reconnect auf der neuen Verbindung neu gesetzt.
- Ein
subscribe, das gesendet wird während der Actor getrennt ist, wird beim nächsten Connect angewandt statt verworfen. { kind: 'unsubscribe', subject }entfernt das Subject endgültig: es kommt beim nächsten Reconnect nicht zurück, auch nicht wenn es aussubscriptionsstammte.- Ein erneutes
subscribeauf einem lebenden Subject tauscht dessen Target — die vorherige Subscription wird zuerst abgebaut, es empfängt also genau ein Actor.
// Zur Laufzeit hinzugefügt; nach einem Broker-Neustart weiterhin da.nats.tell({ kind: 'subscribe', subject: 'audit.>', target: auditor });// Dasselbe Subject auf einen anderen Actor lenken.nats.tell({ kind: 'subscribe', subject: 'audit.>', target: newAuditor });// Endgültig weg.nats.tell({ kind: 'unsubscribe', subject: 'audit.>' });Das betrifft die Subscription, nicht die Nachrichten: Core-NATS hat keine Durability, alles was während der Trennung publiziert wurde ist weiterhin verloren (siehe den Warnhinweis unten).
Subjects (NATS-Topics)
Abschnitt betitelt „Subjects (NATS-Topics)“NATS nutzt mit . getrennte Subjects:
events.user.signupevents.user.deleteorders.priority.placedmetrics.gauge.cpuSubjects sind case-sensitive; per Konvention kleinbuchstabig + punktgetrennt.
Wildcards
Abschnitt betitelt „Wildcards“| Wildcard | Trifft |
|---|---|
* | Genau ein Token (trifft user in events.user.signup). |
> | Ein oder mehr Tokens (trifft user.signup in events.user.signup). |
'events.>' → events.user.signup, events.user.delete, events.order.placed'events.*.signup' → events.user.signup, events.admin.signup> darf nur als letztes Token vorkommen.
Request / Reply
Abschnitt betitelt „Request / Reply“Core-NATS modelliert Request/Reply über ein Reply-Subject: Der
Requester publisht mit einem replyTo-Subject und hört darauf; der
Responder schickt seine Antwort dorthin. Es gibt kein dediziertes
request-Command — du setzt es aus publish / subscribe und dem
replyTo-Feld zusammen.
// Requester — auf einem privaten Reply-Subject hören, dann den// Request mit `replyTo` darauf zeigend publishen:nats.tell({ kind: 'subscribe', subject: 'reply.balance.42', target: replyHandler });nats.tell({ kind: 'publish', publish: { subject: 'account.balance', payload: JSON.stringify({ accountId: '42' }), replyTo: 'reply.balance.42', // Responder antwortet auf diesem Subject },});
// Responder — das Request-Subject subscriben, auf `replyTo` antworten:class BalanceService extends Actor<NatsMessage> { constructor(private readonly nats: ActorRef<NatsCommand>) { super(); } override onReceive(message: NatsMessage): void { if (!message.replyTo) return; this.nats.tell({ kind: 'publish', publish: { subject: message.replyTo, payload: JSON.stringify({ balance: 100 }), }, }); }}Eingehende NatsMessages tragen subject, payload (ein
Uint8Array) und replyTo (ein leerer String, wenn der Publisher
keins gesetzt hat). Das ist das “RPC über NATS”-Muster — ein paar
hundert Mikrosekunden Round-Trip auf localhost.
Wann NATS
Abschnitt betitelt „Wann NATS“Drei primäre Anwendungsfälle:
- Hochvolumiges Pub/Sub ohne die operative Komplexität von Kafka.
- Microservice-Request/Reply — synchrone Calls zwischen Services über Subjects.
- Leichtgewichtige Event-Verteilung — Fire-and-Forget- Benachrichtigungen, Metriken, Log-Streams.
Für durable Streams (Replay, Historie, ACK-Semantik) siehe JetStream unten. Für clusterinternes Pub/Sub ist DistributedPubSub einfacher.
JetStream
Abschnitt betitelt „JetStream“Für NATS mit Durability nimm JetStreamActor:
import { JetStreamActor, JetStreamOptions } from 'actor-ts';
const jetStreamOptions = JetStreamOptions.create() .withServers(['nats://nats-1:4222']) .withStream({ name: 'ORDERS', subjects: ['orders.>'], storage: 'file', maxAge: 86_400 * 7 * 1_000_000_000, // 7 Tage (maxAge ist in Nanosekunden) });const js = system.spawn( () => new JetStreamActor( jetStreamOptions, ), 'js',);JetStream-Schichten:
- Stream — durables Log, Subjects → Records.
- Consumer — Lese-Cursor; mehrere Consumer können denselben Stream unabhängig lesen.
- ACK-Semantik — wie bei AMQP acken Consumer jede Nachricht.
Nimm JetStream, wenn du brauchst:
- Replay — Consumer können zu beliebigen Offsets zurückspulen.
- Persistenz — überlebt Broker-Neustarts (dateibasierter Storage).
- Geordneter Konsum pro Subject.
Kafka schlägt JetStream weiterhin bei massiver Skalierung (Milliarden Events/s über Hunderte Partitionen). JetStream ist der Sweet Spot für “weniger als Kafka, mehr als Core-NATS”.
Peer-Dependency
Abschnitt betitelt „Peer-Dependency“npm install nats# oder: bun add natsDas Paket nats enthält sowohl Core-NATS- als auch JetStream-
Clients.
Wohin als Nächstes
Abschnitt betitelt „Wohin als Nächstes“- I/O-Übersicht — das große Bild.
- BrokerActor-Basis — der gemeinsame Lifecycle.
- Kafka — schwereres durables Streaming.
- MQTT — IoT-fokussierte Alternative.
- DistributedPubSub — für clusterinternes Pub/Sub.
