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, }); } }}Einstellungen
Abschnitt betitelt „Einstellungen“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.
Stream-Konfiguration
Abschnitt betitelt „Stream-Konfiguration“Setze stream, wenn der Actor den Stream beim Verbinden anlegen oder
aktualisieren soll (create ist standardmäßig true).
| Feld | Bedeutung |
|---|---|
name | Stream-Name (erforderlich). |
subjects | Vom Stream erfasste Subjects, z. B. ['orders.>'] (erforderlich). |
retention | 'limits' (Default), 'interest' oder 'workqueue'. |
storage | 'file' (dauerhaft) oder 'memory'. |
maxMessages / maxBytes | Aufbewahrungs-Obergrenzen. |
maxAge | Max. Alter in Nanosekunden (unverändert durchgereicht). |
create | Stream beim Verbinden anlegen/aktualisieren. Default true. |
Consumer-Konfiguration
Abschnitt betitelt „Consumer-Konfiguration“Setze consumer, um eine dauerhafte Subscription zu binden.
| Feld | Bedeutung |
|---|---|
durable | Dauerhafter 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'. |
ackWaitMs | Zeit, bevor der Server ohne Ack erneut zustellt. Default 30_000. |
filterSubject | Subject-Filter — standardmäßig alle Subjects des Streams. |
maxAcknowledgmentPending | Max. unbestätigte Nachrichten in Flight. Default 1024. |
create | Consumer beim Verbinden anlegen/aktualisieren. Default true. |
Befehle
Abschnitt betitelt „Befehle“JetStreamActor akzeptiert einen JetStreamCommand (diskriminiert über kind):
kind | Felder | Zweck |
|---|---|---|
publish | publish: JetStreamPublish | Nachricht publizieren (siehe idempotentes Publish). |
acknowledgment | streamSeq | Zugestellte Nachricht bestätigen — Server markiert sie als konsumiert. |
negativeAcknowledgment | streamSeq, delayMs? | Erneut zustellen (optional nach delayMs). |
terminate | streamSeq, reason? | Endgültiger Fehler — Server verwirft die Nachricht dauerhaft. |
inProgress | streamSeq | Heartbeat — verlängert das Ack-Fenster für einen langen Handler. |
fetch | batch, 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.
Der Ack-Handshake
Abschnitt betitelt „Der Ack-Handshake“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 (nachdelayMs, 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- vs. Pull-Consumer
Abschnitt betitelt „Push- vs. Pull-Consumer“- Push (Default) — der Server streamt Nachrichten; die interne
Pumpe des Actors verteilt jede an
targetund 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 zubatchNachrichten (kehrt nachexpiresMsfrüher zurück, Default5_000); jede Nachricht durchläuft weiterhin denselben Ack-Handshake. Passt besser zu langsamen/schubweisen Consumern als Push-Fan-out.
Idempotentes Publish
Abschnitt betitelt „Idempotentes Publish“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) },});Peer-Dependency
Abschnitt betitelt „Peer-Dependency“npm install nats# oder: bun add natsDas nats-Paket enthält sowohl Core-NATS als auch den JetStream-Client.
Wie geht es weiter
Abschnitt betitelt „Wie geht es weiter“- JetStream KV + Object Store — die Bucket-Sub-APIs.
- NATS — Core-(nicht-dauerhaftes)-Pub/Sub + Request-Reply.
- Kafka — schwereres dauerhaftes Streaming in großem Maßstab.
- BrokerActor-Basis — der gemeinsame Reconnect-/Buffer-/Circuit-Breaker-Lebenszyklus.
- I/O-Überblick — das große Ganze.
