Redis Streams
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' }, },});Settings
Section titled “Settings”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}Consuming with a consumer group
Section titled “Consuming with a consumer group”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
XACKafter 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.
Inbound messages
Section titled “Inbound messages”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.
Pending Entries List (PEL)
Section titled “Pending Entries List (PEL)”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).
Connection loss and reconnect
Section titled “Connection loss and reconnect”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
BrokerConnectedfor a broker nothing has reached. - A dropped connection is reported. The driver’s
error/close/endsignals land onBrokerDisconnected, so a consume-only actor — one that never publishes — is as visible to a health check as a producing one. - A failed
XREADGROUPis 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 … MKSTREAMruns on every connect, so a Redis restart that lost the group heals on the next reconnect instead of spinning onNOGROUP. This is also why aNOGROUPreply 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.
Peer dependency
Section titled “Peer dependency”npm install ioredis# or: bun add ioredisioredis is the underlying client. Versions 5+ are tested.
When to use Redis Streams
Section titled “When to use Redis Streams”Three good fits:
- Existing Redis infrastructure — adding streaming without deploying a new broker.
- Sub-millisecond latency requirements — Redis is fast.
- 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.
Where to next
Section titled “Where to next”- I/O overview — the bigger picture.
- Kafka — when scale exceeds Redis.
- NATS JetStream — durable streams without Redis.
- BrokerActor base — the shared lifecycle.
