Redis Streams
RedisStreamsActor integriert mit Redis Streams (der Familie
XADD / XREADGROUP). Die richtige Form, wenn Redis
bereits in deiner Infrastruktur steckt und du Streaming willst,
ohne Kafka oder NATS aufzusetzen.
import { ActorSystem } from 'actor-ts';import { RedisStreamsActor, RedisStreamsOptions } from 'actor-ts/io';
const redisStreamsOptions = RedisStreamsOptions.create() .withUrl('redis://localhost:6379') .withStreams(['orders', 'payments']) // Consumer-Group-Verdrahtung + Consumer-Actor-Target leben in // den Settings — es gibt kein Runtime-`subscribe`-Surface. .withConsumerGroup({ group: 'order-processors', consumer: 'worker-1' }) .withTarget(orderHandler);const redis = system.spawn( () => new RedisStreamsActor( redisStreamsOptions, ), 'redis',);
// Publishen (XADD unter der Haube):redis.tell({ kind: 'publish', publish: { stream: 'orders', fields: { sku: 'book-1', quantity: '2' }, },});Settings
Abschnitt betitelt „Settings“interface RedisStreamsOptionsType extends BrokerCommonOptionsType { url: string; // 'redis://host:6379' streams?: ReadonlyArray<string>; // Streams, die konsumiert werden consumerGroup?: { group: string; consumer: string; createIfMissing?: boolean; // Default true }; blockMs?: number; // XREADGROUP-Block; Default 5_000 target?: ActorRef<RedisStreamEntry>; // Consumer-Actor für eingehende Einträge tls?: TlsTransportOptionsType; // Cert-Material; aktiviert TLS; nur im Code}Konsumieren mit einer Consumer Group
Abschnitt betitelt „Konsumieren mit einer Consumer Group“Konsum läuft immer über eine Consumer Group — der Actor nutzt
ausschließlich XREADGROUP. Der Consumer startet nur, wenn
consumerGroup, streams und target alle gesetzt sind; es gibt
keinen gruppenlosen XREAD-Broadcast-Pfad.
const redisStreamsOptions = RedisStreamsOptions.create() .withUrl(url) .withStreams(['orders']) .withConsumerGroup({ group: 'workers', consumer: 'worker-1' }) .withTarget(handler);new RedisStreamsActor( redisStreamsOptions,);Consumer-Group-Modus:
- Work-Sharing — jede Nachricht geht an einen Consumer in der Gruppe.
- ACK-Semantik — Consumer rufen nach der Verarbeitung
XACKauf; ungeackte Nachrichten werden in der Pending Entries List (PEL) geführt und können erneut zugestellt werden. - Mehrere Worker — N Actor-Instanzen in derselben Gruppe teilen sich die Last.
Das ist der richtige Modus für Work Queues — Orders gehen an einen Worker, werden verarbeitet, geackt, fertig.
Eingehende Nachrichten
Abschnitt betitelt „Eingehende Nachrichten“class OrderProcessor extends Actor<RedisStreamEntry> { constructor(private readonly redis: ActorRef<RedisStreamsCommand>) { super(); }
override async onReceive(entry: RedisStreamEntry): Promise<void> { const order = JSON.parse(entry.fields.body); await processOrder(order); this.redis.tell({ kind: 'acknowledgment', stream: entry.stream, id: entry.id }); }}Der entry:
type RedisStreamEntry = { stream: string; id: string; // z. B. "1684923847-0" fields: Record<string, string>; // dein XADD-Payload};Stream-IDs sind <millis>-<seq>-Strings. Nimm die id als
Dedup-Key für idempotente Verarbeitung — selbst wenn Redis
erneut zustellt, kann dein Handler schon verarbeitete IDs
überspringen.
Pending Entries List (PEL)
Abschnitt betitelt „Pending Entries List (PEL)“Ungeackte Nachrichten bleiben in der PEL der Consumer Group —
Redis trackt jeden ausgelieferten-aber-nicht-geackten Eintrag pro
Consumer. Nach einem Consumer-Crash können diese Einträge
manuell per XCLAIM gegen den Redis-Server re-claimed werden
(außerhalb des Actor-Surface). Ein zukünftiges Actor-Command
würde das einkapseln; aktuell führe XCLAIM aus einem separaten
ioredis-Client oder Sidecar-Prozess aus.
Für At-Least-Once-Delivery im Normalfall: ack nach der
Verarbeitung tellen — ein Consumer, der mitten in der
Verarbeitung crasht, lässt den Eintrag einfach in der PEL stehen.
Der nächste Consumer, der mit derselben Group startet, kann dort
weitermachen, wo der vorige aufgehört hat (via XPENDING + XCLAIM).
Verbindungsverlust und Reconnect
Abschnitt betitelt „Verbindungsverlust und Reconnect“Der Actor öffnet zwei Verbindungen — einen Producer, einen Consumer —
und führt beide auf den gemeinsamen
BrokerActor-Lifecycle zurück. reconnect
und circuitBreaker bedeuten hier also dasselbe wie bei jedem anderen
Broker:
const redisStreamsOptions = RedisStreamsOptions.create() .withUrl('redis://localhost:6379') .withStreams(['orders']) .withConsumerGroup({ group: 'workers', consumer: 'worker-1' }) .withTarget(handler) .withReconnect({ initialDelayMs: 250, maxDelayMs: 30_000, factor: 2 }) .withCircuitBreaker(5, 60_000);Was das konkret bringt:
- Eine abgelehnte Verbindung lässt den Connect scheitern. Die
Clients werden lazy geöffnet und der Handshake wird abgewartet — ein
nicht erreichbares Redis erzeugt also einen fehlgeschlagenen Versuch
und einen Backoff, kein
BrokerConnectedfür einen Broker, den nie jemand erreicht hat. - Eine abgerissene Verbindung wird gemeldet. Die
error- /close- /end-Signale des Treibers landen aufBrokerDisconnected, sodass ein reiner Consumer-Actor — einer, der nie publisht — für einen Health-Check genauso sichtbar ist wie ein produzierender. - Ein fehlgeschlagenes
XREADGROUPwird klassifiziert. Eine Ablehnung auf Verbindungsebene (Socket geschlossen,ECONNREFUSED, ioredis hat seine Retries pro Request aufgebraucht) beendet die Consumer-Schleife und übergibt den Ausfall an den konfigurierten Backoff. Eine Ablehnung auf Kommandoebene wird an Ort und Stelle wiederholt, und Wiederholungen desselben Fehlers werden zu einem einzigen Log-Eintrag zusammengefasst, der die Anzahl der unterdrückten Vorkommen mitführt. - Die Consumer Group wird neu angelegt.
XGROUP CREATE … MKSTREAMläuft bei jedem Connect — ein Redis-Neustart, der die Group verloren hat, heilt also beim nächsten Reconnect, statt sich aufNOGROUPfestzudrehen. Genau deshalb gilt eineNOGROUP-Antwort als Verbindungsebene: nichts anderes führt das Bootstrapping erneut aus.
ioredis reconnected zusätzlich intern, und die beiden Ebenen
behindern sich nicht: seine Retries sind das, worauf eine Ablehnung auf
Kommandoebene wartet, und die des Frameworks übernehmen, sobald ioredis
aufgibt.
Was es nicht bringt, ist Redelivery: Einträge, die beim
Verbindungsabriss schon an target übergeben, aber noch nicht geackt
waren, bleiben in der PEL, und der Reconnect spielt sie nicht erneut
ein. Die Wiederherstellung läuft über XPENDING + XCLAIM, genau wie
bei einem Consumer-Crash.
Eine rediss://-URL (beachte das doppelte s) startet den Handshake,
und ioredis verifiziert den Server gegen den System-Trust-Store.
Hinter einer privaten CA oder für Client-Zertifikats-Auth kommt
withTls dazu:
const redisStreamsOptions = RedisStreamsOptions.create() .withUrl('rediss://redis.example.com:6380') .withTls({ ca: fs.readFileSync('./tls/ca.crt', 'utf8'), cert: fs.readFileSync('./tls/client.crt', 'utf8'), key: fs.readFileSync('./tls/client.key', 'utf8'), });new RedisStreamsActor(redisStreamsOptions);Jedes Feld trägt das Material selbst, nie einen Pfad darauf.
cert und key ergeben nur gemeinsam Sinn; eines ohne das andere
wird beim Start des Actors abgelehnt.
Producer- und Consumer-Client werden aus demselben Material gebaut —
sie wählen denselben Server, brauchen also dasselbe Vertrauen.
ioredis liest das Vorhandensein der Option als „TLS aushandeln“,
withTls funktioniert also auch gegen eine schlichte
redis://-URL; rediss:// sagt dasselbe nur deutlicher. Einen
Config-Key hat es bewusst nicht: Ein Private Key gehört nicht in die
application.conf.
Peer-Dependency
Abschnitt betitelt „Peer-Dependency“npm install ioredis# oder: bun add ioredisioredis ist der zugrunde liegende Client. Versionen 5+ sind
getestet.
Wann Redis Streams
Abschnitt betitelt „Wann Redis Streams“Drei gute Einsatzfälle:
- Bestehende Redis-Infrastruktur — Streaming hinzufügen, ohne einen neuen Broker auszurollen.
- Anforderungen an Sub-Millisekunden-Latenz — Redis ist schnell.
- Kleinere Workloads — Millionen pro Tag, nicht Milliarden. Für riesige Skalierung ist Kafka zweckgebaut.
Nicht die richtige Form für:
- Langzeit-Retention — Redis Streams können persistieren, sind aber nicht für jahrelange Retention ausgelegt.
- Massive Parallelität — Redis ist single-threaded; ein Redis-Server wird zum Bottleneck.
Wohin als Nächstes
Abschnitt betitelt „Wohin als Nächstes“- I/O-Übersicht — das große Bild.
- Kafka — wenn die Skalierung Redis übersteigt.
- NATS JetStream — durable Streams ohne Redis.
- BrokerActor-Basis — der gemeinsame Lifecycle.
