Zum Inhalt springen
Deutsch

JetStream KV + Object Store

JetStream sind drei APIs, nicht eine. JetStreamActor deckt Streams und Consumer ab; diese Seite behandelt die beiden anderen:

  • JetStreamKeyValueActor — ein KV-Bucket: revisionierte Keys, Compare-and-Swap-Schreibvorgänge und ein watch-Change-Feed.
  • JetStreamObjectStoreActor — ein Object-Bucket: benannte Blobs mit Metadaten.

Beide sind eigene Actors, weil es eigene Sub-APIs mit eigener Semantik sind — ein Bucket ist kein Stream + Consumer, und sie in den Stream-Actor zu falten hätte ihm einen sechsten Modus verpasst, den niemand will.

import { ActorSystem, Actor, JetStreamKeyValueActor, JetStreamKeyValueOptions } from 'actor-ts';
import type { ActorRef, JetStreamKeyValueMessage } from 'actor-ts';
const sessionOptions = JetStreamKeyValueOptions.create()
.withServers(['nats://localhost:4222'])
.withBucket('sessions')
.withHistory(5)
.withTimeToLive(3_600_000);
const sessions = system.spawn(() => new JetStreamKeyValueActor(sessionOptions), 'sessions');
// Fire-and-forget write:
sessions.tell({ kind: 'put', key: 'user.42', value: JSON.stringify(session) });
// Read — the answer arrives at `reader` as a kind-tagged message:
sessions.tell({ kind: 'get', key: 'user.42', target: reader });

Der Builder ist der primäre Stil; ein einfaches Objekt geht ebenso. Die gemeinsamen Broker-Felder — withReconnect, withCircuitBreaker, withOutboundBuffer — kommen von der gemeinsamen BrokerActor-Basis. servers und bucket sind Pflicht.

kindFelderZweck
putkey, value, expectedRevision?, target?Wert schreiben; mit expectedRevision als Compare-and-Swap.
getkey, targetAktuellen Wert lesen.
deletekey, target?Key entfernen; ein Tombstone bleibt in der Historie.
purgekey, target?Key und seine Historie entfernen.
keysfilter?, targetLebende Keys auflisten, optional per Subject-Filter eingegrenzt.
watchkey?, targetJede Änderung unter key streamen (Default '>' — das ganze Bucket).
unwatchkey?Diesen Watch beenden.

target ist Pflicht, wo die Antwort der Zweck ist (get, keys, watch), und optional, wo sie eine Quittung ist (put, delete, purge) — lass es für Fire-and-forget-Schreibvorgänge weg, dann werden Fehler stattdessen geloggt.

Jede Antwort ist eine JetStreamKeyValueMessage, diskriminiert über kind:

kindFelderWann
keyValueEntrykey, value, revision, createdAtEin get-Treffer und jedes watch-Update.
keyValueNotFoundkeyEin get auf einen fehlenden oder gelöschten Key.
keyValueRemovedkey, purgedQuittung für delete / purge und jede beobachtete Entfernung.
keyValueRevisionkey, revisionQuittung für ein erfolgreiches put.
keyValueKeyskeysAntwort auf keys.
keyValueOperationFailedoperation, key?, reasonDie Operation ist fehlgeschlagen — meist ein Compare-and-Swap-Konflikt.
class SessionReader extends Actor<JetStreamKeyValueMessage> {
override onReceive(message: JetStreamKeyValueMessage): void {
if (message.kind === 'keyValueEntry') {
this.log.info(`${message.key} @ rev ${message.revision}`);
}
}
}

revision ist das Concurrency-Token. Lies es, berechne den neuen Wert und schreibe ihn mit expectedRevision zurück — der Server weist den Schreibvorgang ab, falls jemand anderes den Key zwischenzeitlich bewegt hat, und der Actor antwortet mit keyValueOperationFailed, sodass der Aufrufer gegen den frischen Wert erneut versuchen kann:

// after a `get` answered with { revision: 7 }
sessions.tell({
kind: 'put',
key: 'user.42',
value: JSON.stringify(next),
expectedRevision: 7,
target: reader,
});

expectedRevision: 0 bedeutet „nur, wenn der Key noch nicht existiert”.

Ein watch ist Soll-Zustand, kein einmaliges Abonnement: Er wird bei jedem Reconnect wiederhergestellt, und einer, der während einer Trennung abgesetzt wird, landet beim nächsten Connect, statt verworfen zu werden. Das ist dieselbe Zusage, die NatsActor einer Subscription gibt — und sie wiegt hier schwerer: Ein Change-Feed, der still stehenbleibt, sieht genauso aus wie ein Bucket, das sich nicht mehr ändert.

sessions.tell({ kind: 'watch', key: 'user.>', target: reader });
// … later
sessions.tell({ kind: 'unwatch', key: 'user.>' });

Ein erneutes watch auf einen laufenden Key tauscht das Target aus.

interface JetStreamKeyValueOptionsType extends BrokerCommonOptionsType {
servers?: string[] | string; // NATS server URLs (required)
token?: string;
user?: string;
password?: string;
name?: string; // client identifier
bucket?: string; // bucket name (required)
history?: number; // revisions kept per key (create-time)
timeToLive?: number; // per-key TTL in ms (create-time)
storage?: 'memory' | 'file'; // create-time
replicas?: number; // create-time
maxValueBytes?: number; // server-side cap on one value (create-time)
create?: boolean; // create the bucket when missing; default true
}

Die Create-Time-Felder werden nur gesendet, wenn der Actor das Bucket anlegt. Mit create: false bindet der Actor an ein vorhandenes Bucket und lässt den Connect scheitern, wenn es fehlt — die richtige Einstellung, wenn der Betrieb das Bucket bereitstellt und ein Tippfehler nicht stillschweigend ein leeres erzeugen soll.

HOCON-Defaults liegen unter actor-ts.io.broker.jetstream-key-value:

actor-ts.io.broker.jetstream-key-value {
servers = ["nats://localhost:4222"]
bucket = "sessions"
history = 5
}
import { JetStreamObjectStoreActor, JetStreamObjectStoreOptions } from 'actor-ts';
const assetOptions = JetStreamObjectStoreOptions.create()
.withServers(['nats://localhost:4222'])
.withBucket('assets');
const assets = system.spawn(() => new JetStreamObjectStoreActor(assetOptions), 'assets');
assets.tell({ kind: 'put', name: 'report.pdf', payload: bytes, target: uploader });
assets.tell({ kind: 'get', name: 'report.pdf', target: reader });
kindFelderZweck
putname, payload, description?, headers?, target?Body speichern und einen vorherigen ersetzen.
getname, targetDen ganzen Body lesen.
deletename, target?Objekt entfernen.
infoname, targetNur Metadaten — keine Body-Übertragung, unabhängig von der Größe.
listtargetMetadaten aller lebenden Objekte im Bucket.
kindFelderWann
objectStoredname, infoQuittung für put.
objectBodyname, payload, infoAntwort auf get.
objectInfoname, infoAntwort auf info.
objectListobjectsAntwort auf list.
objectDeletednameQuittung für delete.
objectNotFoundnameDas Objekt fehlt oder ist gelöscht.
objectStoreOperationFailedoperation, name?, reasonDie Operation ist fehlgeschlagen, auch bei einem Body über der Grenze.

info trägt name, size, chunks, digest, modifiedAt sowie die optionalen description / headers, mit denen das Objekt gespeichert wurde.

interface JetStreamObjectStoreOptionsType extends BrokerCommonOptionsType {
servers?: string[] | string; // NATS server URLs (required)
token?: string;
user?: string;
password?: string;
name?: string; // client identifier
bucket?: string; // bucket name (required)
description?: string; // create-time
storage?: 'memory' | 'file'; // create-time
replicas?: number; // create-time
maxObjectBytes?: number; // whole-object ceiling; default 1 MiB
create?: boolean; // create the bucket when missing; default true
}

HOCON-Defaults liegen unter actor-ts.io.broker.jetstream-object-store:

actor-ts.io.broker.jetstream-object-store {
servers = ["nats://localhost:4222"]
bucket = "assets"
maxObjectBytes = 4M
}

Erhöhe maxObjectBytes nur so weit, wie outboundBuffer × maxObjectBytes ein Speicherwert bleibt, den du über einen Ausfall hinweg halten willst; setze outboundBuffer auf 0, wenn dir Fail-fast lieber ist als Blobs zu puffern. info und list sind davon nicht betroffen — Metadaten tragen keinen Body und antworten für Objekte jeder Größe.

Für wirklich große Objekte nutze die gechunkten put/get-Streams des nats-Clients direkt. Sie durch den Actor zu führen braucht einen Dispatch-Pfad außerhalb der Verbindungs-Zustandsmaschine — eine größere Änderung als dieser Actor.

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

Dasselbe nats-Paket, das NATS und JetStream trägt — die KV- und Object-Store-Views kommen mit.

  • NATS JetStream — dauerhafte Streams und Consumer, die dritte Sub-API.
  • NATS — Core-(nicht-dauerhaftes)-Pub/Sub + Request-Reply.
  • BrokerActor-Basis — der gemeinsame Reconnect-/Buffer-/Circuit-Breaker-Lebenszyklus.
  • I/O-Überblick — das große Ganze.