Aller au contenu
Français

Redis Streams

Ce contenu n’est pas encore disponible dans votre langue.

RedisStreamsActor integrates with Redis Streams (the XADD / XREADGROUP family). Right shape when Redis is already in your infrastructure and you want streaming without deploying Kafka or NATS.

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 wiring + consumer-actor target live in
// settings — there is no runtime `subscribe` surface.
.withConsumerGroup({ group: 'order-processors', consumer: 'worker-1' })
.withTarget(orderHandler);
const redis = system.spawn(
() => new RedisStreamsActor(
redisStreamsOptions,
),
'redis',
);
// Publish (XADD under the hood):
redis.tell({
kind: 'publish',
publish: {
stream: 'orders',
fields: { sku: 'book-1', quantity: '2' },
},
});
interface RedisStreamsOptionsType extends BrokerCommonOptionsType {
url: string; // 'redis://host:6379'
streams?: ReadonlyArray<string>; // streams to consume
consumerGroup?: {
group: string;
consumer: string;
createIfMissing?: boolean; // default true
};
blockMs?: number; // XREADGROUP block; default 5_000
target?: ActorRef<RedisStreamEntry>; // consumer-actor for inbound entries
tls?: TlsTransportOptionsType; // cert material; also enables TLS; code-only
}

Consumption always goes through a consumer group — the actor uses XREADGROUP exclusively. The consumer starts only when consumerGroup, streams, and target are all set; there is no group-less XREAD broadcast path.

const redisStreamsOptions = RedisStreamsOptions.create()
.withUrl(url)
.withStreams(['orders'])
.withConsumerGroup({ group: 'workers', consumer: 'worker-1' })
.withTarget(handler);
new RedisStreamsActor(
redisStreamsOptions,
);

Consumer-group mode:

  • Work-sharing — each message goes to one consumer in the group.
  • Ack semantics — consumers XACK after processing; un-acked messages are tracked in the pending entries list (PEL) and can be re-delivered.
  • Multiple workers — N actor instances in the same group share the workload.

This is the right mode for work queues — orders go to one worker, processed, acked, done.

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 });
}
}

The entry:

type RedisStreamEntry = {
stream: string;
id: string; // e.g. "1684923847-0"
fields: Record<string, string>; // your XADD payload
};

Stream IDs are <millis>-<seq> strings. Use the id as a dedup key for idempotent processing — even if Redis re-delivers, your handler can skip already-processed IDs.

Un-acked messages stay in the consumer group’s PEL — Redis tracks every delivered-but-not-acked entry per consumer. After a consumer crash, those entries can be manually re-claimed via XCLAIM against the Redis server (outside the actor surface). A future actor command would wrap that; for now, run XCLAIM from a separate ioredis client or use a sidecar process.

For at-least-once delivery in the common path: tell ack after processing, and a consumer that crashes mid-processing simply leaves the entry in the PEL — the next consumer that starts with the same group can pick up where the previous one stopped (via XPENDING + XCLAIM).

The actor opens two connections — one producer, one consumer — and routes both onto the shared BrokerActor lifecycle, so reconnect and circuitBreaker mean the same thing here as for every other 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);

What that buys, concretely:

  • A refused connection fails the connect. The clients are opened lazily and the handshake is awaited, so an unreachable Redis produces a failed attempt and a backoff — not a BrokerConnected for a broker nothing has reached.
  • A dropped connection is reported. The driver’s error / close / end signals land on BrokerDisconnected, so a consume-only actor — one that never publishes — is as visible to a health check as a producing one.
  • A failed XREADGROUP is classified. A connection-level rejection (socket closed, ECONNREFUSED, ioredis exhausting its per-request retries) ends the consumer loop and hands the outage to the configured backoff. A command-level rejection retries in place, and repeats of an identical failure are collapsed into one log record carrying the count of what it stood in for.
  • The consumer group is re-created. XGROUP CREATE … MKSTREAM runs on every connect, so a Redis restart that lost the group heals on the next reconnect instead of spinning on NOGROUP. This is also why a NOGROUP reply is treated as connection-level: nothing else re-runs the bootstrap.

ioredis also reconnects underneath, and the two layers do not fight: its retries are what a command-level rejection waits on, and the framework’s take over once it gives up.

What it does not buy is redelivery: entries already handed to target but not yet acked when the connection drops stay in the PEL, and the reconnect does not replay them. Recovery is XPENDING + XCLAIM, exactly as for a consumer crash.

A rediss:// URL (note the double-s) starts the handshake and ioredis verifies the server against the system trust store. Behind a private CA, or for client-certificate auth, add withTls:

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);

Every field carries the material itself, never a path to it. cert and key only mean anything together, and one without the other is rejected when the actor starts.

The producer and the consumer client are built from the same material — they dial the same server, so they need the same trust. ioredis reads the option’s presence as “negotiate TLS”, so withTls also works against a plain redis:// URL; rediss:// is the clearer way to say the same thing. It has no config key on purpose: a private key does not belong in application.conf.

Terminal window
npm install ioredis
# or: bun add ioredis

ioredis is the underlying client. Versions 5+ are tested.

Three good fits:

  1. Existing Redis infrastructure — adding streaming without deploying a new broker.
  2. Sub-millisecond latency requirements — Redis is fast.
  3. Smaller workloads — millions per day, not billions. For huge scale, Kafka is purpose-built.

Not the right shape for:

  • Long-term retention — Redis Streams can persist but aren’t designed for years-long retention.
  • Massive parallelism — Redis is single-threaded; one Redis server bottlenecks.