コンテンツにスキップ
日本語

NATS JetStream

このコンテンツはまだ日本語訳がありません。

JetStreamActor integrates with NATS JetStream — the durable streaming layer on top of core NATS. Where NatsActor is fire-and-forget pub/sub, JetStreamActor gives you persistent streams, replay, and explicit acknowledgment with Kafka-style “exactly-once-with-processing” semantics via the 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 when messageId is set):
jetStream.tell({
kind: 'publish',
publish: {
subject: 'orders.created',
payload: JSON.stringify(order),
messageId: order.id, // server dedupes within the stream's dedup window
},
});

The consumer forwards every message to target and waits for an explicit ack before the ack-window (ackWaitMs) expires — if none arrives, the server redelivers:

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 {
// redeliver after 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 identifier
stream?: JetStreamStreamConfig; // set when this actor owns the stream
consumer?: JetStreamConsumerConfig; // required to start a subscription
target?: ActorRef<JetStreamMessage>; // receives every consumed message
acknowledgmentTimeout?: number; // ack-wait ceiling; default consumer.ackWaitMs ?? 30_000
}

The builder is the primary style; a plain object works too. Common broker fields — withReconnect, withCircuitBreaker, withOutboundBuffer — come from the shared BrokerActor base. servers is required.

Set stream when the actor should create or update the stream at connect time (create defaults to true).

FieldMeaning
nameStream name (required).
subjectsSubjects captured by the stream, e.g. ['orders.>'] (required).
retention'limits' (default), 'interest', or 'workqueue'.
storage'file' (durable) or 'memory'.
maxMessages / maxBytesRetention caps.
maxAgeMax age in nanoseconds (passed through verbatim).
createCreate/update the stream at connect time. Default true.

Set consumer to bind a durable subscription.

FieldMeaning
durableDurable consumer name (required — survives restarts).
mode'push' (default) or 'pull' (see below).
deliverPolicy'all' (default), 'last', 'new', { kind: 'byStartSeq', startSeq }, or { kind: 'byStartTime', startTimeMs }.
ackPolicy'explicit' (default — ack/nak/term required), 'none', or 'all'.
ackWaitMsTime before the server redelivers without an ack. Default 30_000.
filterSubjectSubject filter — defaults to all subjects in the stream.
maxAcknowledgmentPendingMax in-flight unacked messages. Default 1024.
createCreate/update the consumer at connect time. Default true.

JetStreamActor accepts a JetStreamCommand (discriminated on kind):

kindFieldsPurpose
publishpublish: JetStreamPublishPublish a message (see idempotent publish).
acknowledgmentstreamSeqAck a delivered message — the server marks it consumed.
negativeAcknowledgmentstreamSeq, delayMs?Redeliver (optionally after delayMs).
terminatestreamSeq, reason?Terminal failure — server drops the message permanently.
inProgressstreamSeqHeartbeat — extend the ack-wait window for a long handler.
fetchbatch, expiresMs?Pull mode only — request up to batch messages.

Every inbound JetStreamMessage carries subject, payload (Uint8Array), replyTo, streamSeq, consumerSeq, deliveries (delivery count — 1 on first try), timestamp, and headers.

With the default ackPolicy: 'explicit', every delivered message must be resolved with exactly one of acknowledgment, negativeAcknowledgment, or terminate, keyed by streamSeq:

  • acknowledgment — processing succeeded; the server advances the consumer.
  • negativeAcknowledgment — retry; the server redelivers (after delayMs if given).
  • terminate — give up permanently; the server will not redeliver.
  • inProgress — not terminal; resets the ack-wait timer so a slow handler isn’t redelivered mid-flight.

Set ackPolicy: 'none' for pure fire-and-hose delivery where no ack is expected.

  • Push (default) — the server streams messages; the actor’s internal pump fans each one to target and waits per-message for the ack handshake. The natural fit for actor-style fan-out.
  • Pull (mode: 'pull') — the application self-paces by sending { kind: 'fetch', batch, expiresMs }. Each fetch delivers up to batch messages (returning early after expiresMs, default 5_000); every message still goes through the same ack handshake. Fits slow/bursty consumers better than push fan-out.

JetStreamPublish supports JetStream’s server-side dedup and optimistic concurrency:

jetStream.tell({
kind: 'publish',
publish: {
subject: 'orders.created',
payload: JSON.stringify(order),
messageId: order.id, // sent as Nats-Msg-Id — dedupes re-publishes in the stream's window
expectedLastSeq: lastSeq, // server rejects the publish if the stream head moved (optimistic concurrency)
},
});
Terminal window
npm install nats
# or: bun add nats

The nats package includes both core NATS and the JetStream client.