Zum Inhalt springen
Deutsch

NATS JetStream

JetStreamActor integriert NATS JetStream — die dauerhafte Streaming-Schicht über Core-NATS. Wo NatsActor Fire-and-forget-Pub/Sub ist, liefert JetStreamActor persistente Streams, Replay und explizite Bestätigung mit Kafka-artiger „exactly-once-with-processing”-Semantik über den acknowledgment / negativeAcknowledgment / terminate-Handshake.

import { ActorSystem, Actor, JetStreamActor, JetStreamOptions } from 'actor-ts';
import type { ActorRef, JetStreamMessage, JetStreamCommand } from 'actor-ts';
const jetStreamOptions = JetStreamOptions.create()
.withServers(['nats://localhost:4222'])
.withStream({ name: 'ORDERS', subjects: ['orders.>'] })
.withConsumer({ durable: 'order-processor', ackWaitMs: 30_000 })
.withTarget(orderProcessor);
const jetStream = system.spawn(
() => new JetStreamActor(
jetStreamOptions,
),
'orders-stream',
);
// Publish (idempotent, wenn messageId gesetzt ist):
jetStream.tell({
kind: 'publish',
publish: {
subject: 'orders.created',
payload: JSON.stringify(order),
messageId: order.id, // Server dedupliziert innerhalb des Dedup-Fensters des Streams
},
});

Der Consumer leitet jede Nachricht an target weiter und wartet auf eine explizite Bestätigung, bevor das Ack-Fenster (ackWaitMs) abläuft — kommt keine, stellt der Server erneut zu:

class OrderProcessor extends Actor<JetStreamMessage> {
constructor(private readonly jetStream: ActorRef<JetStreamCommand>) { super(); }
async onReceive(message: JetStreamMessage): Promise<void> {
try {
await db.insertOrder(JSON.parse(new TextDecoder().decode(message.payload)));
this.jetStream.tell({ kind: 'acknowledgment', streamSeq: message.streamSeq });
} catch {
// erneute Zustellung nach 5 s
this.jetStream.tell({
kind: 'negativeAcknowledgment',
streamSeq: message.streamSeq,
delayMs: 5_000,
});
}
}
}
interface JetStreamOptionsType extends BrokerCommonOptionsType {
servers?: string[] | string; // NATS-Server-URLs
token?: string;
user?: string;
password?: string;
name?: string; // Client-Bezeichner
stream?: JetStreamStreamConfig; // gesetzt, wenn dieser Actor den Stream besitzt
consumer?: JetStreamConsumerConfig; // nötig, um eine Subscription zu starten
target?: ActorRef<JetStreamMessage>; // empfängt jede konsumierte Nachricht
acknowledgmentTimeout?: number; // Ack-Wait-Obergrenze; Default consumer.ackWaitMs ?? 30_000
}

Der Builder ist der primäre Stil; ein einfaches Objekt funktioniert ebenso. Gemeinsame Broker-Felder — withReconnect, withCircuitBreaker, withOutboundBuffer — stammen aus der BrokerActor-Basis. servers ist erforderlich.

Setze stream, wenn der Actor den Stream beim Verbinden anlegen oder aktualisieren soll (create ist standardmäßig true).

FeldBedeutung
nameStream-Name (erforderlich).
subjectsVom Stream erfasste Subjects, z. B. ['orders.>'] (erforderlich).
retention'limits' (Default), 'interest' oder 'workqueue'.
storage'file' (dauerhaft) oder 'memory'.
maxMessages / maxBytesAufbewahrungs-Obergrenzen.
maxAgeMax. Alter in Nanosekunden (unverändert durchgereicht).
createStream beim Verbinden anlegen/aktualisieren. Default true.

Setze consumer, um eine dauerhafte Subscription zu binden.

FeldBedeutung
durableDauerhafter Consumer-Name (erforderlich — übersteht Neustarts).
mode'push' (Default) oder 'pull' (siehe unten).
deliverPolicy'all' (Default), 'last', 'new', { kind: 'byStartSeq', startSeq } oder { kind: 'byStartTime', startTimeMs }.
ackPolicy'explicit' (Default — ack/nak/term erforderlich), 'none' oder 'all'.
ackWaitMsZeit, bevor der Server ohne Ack erneut zustellt. Default 30_000.
filterSubjectSubject-Filter — standardmäßig alle Subjects des Streams.
maxAcknowledgmentPendingMax. unbestätigte Nachrichten in Flight. Default 1024.
createConsumer beim Verbinden anlegen/aktualisieren. Default true.

JetStreamActor akzeptiert einen JetStreamCommand (diskriminiert über kind):

kindFelderZweck
publishpublish: JetStreamPublishNachricht publizieren (siehe idempotentes Publish).
acknowledgmentstreamSeqZugestellte Nachricht bestätigen — Server markiert sie als konsumiert.
negativeAcknowledgmentstreamSeq, delayMs?Erneut zustellen (optional nach delayMs).
terminatestreamSeq, reason?Endgültiger Fehler — Server verwirft die Nachricht dauerhaft.
inProgressstreamSeqHeartbeat — verlängert das Ack-Fenster für einen langen Handler.
fetchbatch, expiresMs?Nur Pull-Modus — bis zu batch Nachrichten anfordern.

Jede eingehende JetStreamMessage trägt subject, payload (Uint8Array), replyTo, streamSeq, consumerSeq, deliveries (Zustellungszähler — 1 beim ersten Versuch), timestamp und headers.

Mit dem Default ackPolicy: 'explicit' muss jede zugestellte Nachricht mit genau einem von acknowledgment, negativeAcknowledgment oder terminate aufgelöst werden, adressiert über streamSeq:

  • acknowledgment — Verarbeitung erfolgreich; der Server rückt den Consumer vor.
  • negativeAcknowledgment — erneut versuchen; der Server stellt erneut zu (nach delayMs, falls angegeben).
  • terminate — endgültig aufgeben; keine erneute Zustellung.
  • inProgress — nicht terminal; setzt den Ack-Wait-Timer zurück, damit ein langsamer Handler nicht mitten in der Arbeit erneut zugestellt wird.

Setze ackPolicy: 'none' für reine Fire-and-Hose-Zustellung ohne erwartetes Ack.

  • Push (Default) — der Server streamt Nachrichten; die interne Pumpe des Actors verteilt jede an target und wartet pro Nachricht auf den Ack-Handshake. Der natürliche Fit für Actor-artiges Fan-out.
  • Pull (mode: 'pull') — die Anwendung bestimmt das Tempo selbst, indem sie { kind: 'fetch', batch, expiresMs } sendet. Jedes Fetch liefert bis zu batch Nachrichten (kehrt nach expiresMs früher zurück, Default 5_000); jede Nachricht durchläuft weiterhin denselben Ack-Handshake. Passt besser zu langsamen/schubweisen Consumern als Push-Fan-out.

JetStreamPublish unterstützt JetStreams serverseitige Deduplizierung und optimistische Nebenläufigkeit:

jetStream.tell({
kind: 'publish',
publish: {
subject: 'orders.created',
payload: JSON.stringify(order),
messageId: order.id, // als Nats-Msg-Id gesendet — dedupliziert Re-Publishes im Fenster des Streams
expectedLastSeq: lastSeq, // Server lehnt das Publish ab, wenn der Stream-Kopf sich bewegt hat (optimistische Nebenläufigkeit)
},
});
Terminal-Fenster
npm install nats
# oder: bun add nats

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