Zum Inhalt springen
Deutsch

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),
},
});
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 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 aus subscriptions stammte.
  • Ein erneutes subscribe auf 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).

NATS nutzt mit . getrennte Subjects:

events.user.signup
events.user.delete
orders.priority.placed
metrics.gauge.cpu

Subjects sind case-sensitive; per Konvention kleinbuchstabig + punktgetrennt.

WildcardTrifft
*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.

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.

Drei primäre Anwendungsfälle:

  1. Hochvolumiges Pub/Sub ohne die operative Komplexität von Kafka.
  2. Microservice-Request/Reply — synchrone Calls zwischen Services über Subjects.
  3. 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.

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

Terminal-Fenster
npm install nats
# oder: bun add nats

Das Paket nats enthält sowohl Core-NATS- als auch JetStream- Clients.