Redis · Streams

Streams

A stream is an append-only log: every entry has an identifier assigned by the server and a set of fields. It differs from a list in that reading does not take an entry away, and consumer groups make Redis itself keep the books: who took what, and what has not been acknowledged yet.

Package flytachi/winter-redisClasses RedisStream · RedisStreamGroupEntries StreamEntry · PendingEntry

What a stream is and why

The problem. A queue on a list answers one question — “what next” — and loses everything else. What is read disappears: an event cannot be re-read, one event cannot go to two different handlers, and after an incident there is no way to look at what actually happened. And if a worker takes a job and dies, the job goes with it, with nowhere to learn that from: the “taken but not finished” bookkeeping has to be built by hand.

The solution. A stream is a log that reading does not consume. The entry stays where it is, it has an identifier, and that identifier is how you go back or continue from where you stopped. Consumer groups then hand the delivery bookkeeping to Redis itself: it remembers who got which entry and whether they acknowledged it.

php
$events = $store->stream('events');

$events->add(['type' => 'signup', 'user' => '42']);   // '1755600000123-0'

$events->ensureGroup('mailers');
$group = $events->group('mailers', consumer: 'worker-1');

foreach ($group->consume(count: 10, timeout: 5) as $entry) {
  $this->handle($entry);
  $group->ack($entry);        // until acknowledged, the entry is charged to this worker
}

Properties of the structure

An entry looks like this: the identifier 1755600000123-0 (milliseconds of the server clock and a sequence number within that millisecond) and the fields {"type": "signup", "user": "42"}. Identifiers grow monotonically, so “read everything after this one” is an ordinary operation rather than a search.

The properties everything else follows from:

  • Reading takes nothing away. An entry lives until it is deleted or the log is trimmed. The same entry can be read by any number of independent readers.
  • Order is guaranteed and matches the order of appending.
  • The log grows forever unless it is capped. That is the main price of history.
  • The server keeps the group’s books. It remembers which entry went to which consumer and whether it was acknowledged. A list cannot do that — there the bookkeeping is manual.

How this differs from a list

List Stream
After reading the element is gone the entry stays
Readers per message one as many as you like
History none yes, until trimmed
“Taken but unfinished” tracking manual (moveTo + remove) built in, per group
Growth bounded by consumption forever, unless capped
Identifiers none yes; used for reading and acknowledging

If you need a queue and no history, take a list: it is simpler and cheaper. A stream is for history, several independent readers, or built-in delivery tracking.

Where it is normally used

An event bus. One event, several independent handlers: an email, a metric, a webhook. Each reads with its own group and does not get in the others’ way.

A queue with a guarantee. A worker took an entry and crashed — it stays in the pending list and goes to another worker. Exactly what has to be assembled from moveTo() and remove() with lists.

A change log. An audit trail, an order’s history, an activity feed: it can be re-read, rewound, and an incident can be reconstructed from it.

A buffer with catch-up reading. A handler writes, a consumer reads at its own pace and after a restart continues where it stopped.

Where a stream is the wrong choice:

Task Why not a stream What instead
A simple queue with no history needless growth and needless concepts a list
Notify everyone connected right now a stream stores, it does not broadcast pub/sub
Find an entry by its contents there is no search, only id and range an index on sets
Store an object’s state a stream is about events, not the current value a hash

Entries: StreamEntry

The reading methods return StreamEntry objects rather than raw arrays.

The driver hands entries over as ['1755600000123-0' => ['type' => 'signup']] — an array whose single key is the identifier. Walking that means key()/current(), while the identifier is needed constantly: it is what acknowledges work and what remembers a position. StreamEntry puts it next to the fields.

php
foreach ($events->range(count: 100) as $entry) {
  $entry->id;                  // '1755600000123-0'
  $entry->fields;              // ['type' => 'signup', 'user' => '42']
  $entry->get('type');         // 'signup'
  $entry->get('absent', '—');  // a default value
  $entry->has('user');         // bool
  $entry->timestamp();         // 1755600000123 — milliseconds of the server clock
}

RedisStream reference

RedisStream is a handle on one stream: the object through which commands against that key are run. It is taken from the store:

php
$events = $store->stream('events');    // the server will see 'session:events'

stream() executes nothing — it is a view, not a request. The prefix is applied once, here. The handle holds no connection, with one exception: follow() opens one of its own, because a blocking read occupies a connection for the whole wait.

Writing

add()

Appends an entry to the end of the log.

Syntax

php
public function add(
  array $fields,
  ?int $cap = null,
  bool $exact = false,
  string $id = '*',
): string

Parameters

$fields — the entry’s fields. At least one: Redis has no empty entries. Values are strings, as everywhere.

$cap — cap the log at roughly this many entries as part of the append. null by default — no cap.

$exact — trim exactly to $cap rather than approximately. false by default; the difference is covered under trim().

$id — the entry’s identifier. '*' by default — the server assigns it. A custom identifier must be greater than the last one in the log.

Returns

The identifier assigned to the entry.

Errors

RedisCommandException — if the key is occupied by a structure of another type, or if the given $id is not greater than the last one:

text
RedisCommandException: ERR The ID specified in XADD is equal or smaller
than the target stream top item

Example

php
$events->add(['type' => 'signup', 'user' => '42']);
// '1755600000123-0'

$events->add(['level' => 'warn', 'msg' => $text], cap: 10_000);

Custom identifiers are almost never needed: the server’s are monotonic and carry the time. Assigning them by hand only makes sense when migrating data from another log.

Field values are strings

As everywhere in Redis. An array in a field turns into "Array" unless the configuration has a serializer.

trim()

Keeps roughly $cap of the newest entries.

Syntax

php
public function trim(int $cap, bool $exact = false): int

Parameters

$cap — how many entries to keep.

$exact — trim exactly rather than approximately. false by default.

Returns

The number of entries deleted.

Example

php
$events->trim(10_000);                 // cheap, the length is "about that"
$events->trim(10_000, exact: true);    // exactly that many, more expensive

An approximate trim may delete nothing at all

It stops at the boundary of the stream’s internal node — that is cheaper, and it is what Redis recommends. The consequence is not obvious, and it was verified against a live server: on a log of ten entries trim(5) returns 0 and leaves all ten, while trim(5, exact: true) deletes five and leaves five.

This is not a failure. If you need a predictable length, ask for exact: true and pay for the walk.

trimBefore()

Deletes everything older than the given identifier.

This is how a log is capped by time rather than by count: an identifier starts with the milliseconds of the server clock, so “older than a day” is simply an identifier.

Syntax

php
public function trimBefore(string $id, bool $exact = false): int

Parameters

$id — the boundary; entries with a smaller identifier are deleted.

$exact — an exact boundary instead of an approximate one. false by default, and the caveat from trim() applies here in exactly the same way.

Returns

The number of entries deleted.

Example

php
$events->trimBefore((string) ((time() - 86400) * 1000));

delete()

Deletes entries by identifier.

Deleting from the middle of a log is allowed, but it is the exception: the usual way to bound a stream is trimming. delete() is for when one particular entry must not be shown any more.

Syntax

php
public function delete(string ...$ids): int

Parameters

...$ids — entry identifiers, as many as you like. You may pass none.

Returns

The number of entries that existed and were deleted. With no arguments — 0.

Reading

count()

Reports the number of entries in the log.

Syntax

php
public function count(): int

Parameters

None.

Returns

The number of entries; 0 for a key that does not exist.

range()

Returns entries in append order, taking nothing away.

Syntax

php
public function range(?string $from = null, ?string $to = null, ?int $count = null): array

Parameters

$from — the identifier to read from, inclusive. null by default — from the start of the log.

$to — the identifier to read up to, inclusive. null by default — to the end.

$count — no more than this many entries. null by default — no limit.

Returns

An array of StreamEntry; an empty array for an empty log.

Example

php
$events->range();                    // everything
$events->range(count: 100);          // the first hundred
$events->range($since, count: 50);   // from a known position

reverse()

Returns entries newest first.

A cheap way to see what just happened without reading the log from the start.

Syntax

php
public function reverse(?int $count = null): array

Parameters

$count — no more than this many entries. null by default — the whole log.

Returns

An array of StreamEntry, freshest first.

Example

php
$events->reverse(count: 10);   // the last ten, newest on top

after()

Returns entries appended after the given identifier, without waiting for new ones.

This is the basis of a reader that keeps its own position: remember the identifier of the last entry you handled and pass it in next time. Two such readers do not interfere — a stream is not consumed by reading.

Syntax

php
public function after(string $id, int $count = 10): array

Parameters

$id — the position to read after. The entry with that identifier is not included.

$count — no more than this many entries. 10 by default.

Returns

An array of StreamEntry; empty if there is nothing new.

Example

php
$position = $this->positions->get('events') ?? '0';

foreach ($events->after($position, count: 100) as $entry) {
  $this->handle($entry);
  $position = $entry->id;
}

$this->positions->set('events', $position);

follow()

Waits for entries newer than the ones this handle has already seen.

Syntax

php
public function follow(float $timeout = 0.0, int $count = 10, ?string $from = null): array

Parameters

$timeout — how many seconds to wait. 0.0 by default — indefinitely.

$count — no more than this many entries per call. 10 by default.

$from — the position for the first call. null by default — start at the current end of the log; '0' — read the whole history first.

Returns

An array of StreamEntry; empty if nothing appeared in the time allowed.

Example

main/Daemons/EventTail.php
$events = $store->stream('events');

while (!$this->stopping) {
  foreach ($events->follow(timeout: 5) as $entry) {
      $this->handle($entry);
  }
}

$events->close();

The position is remembered — and why it is not `$`

The first call starts at the current end of the log and spends one extra request on that: it asks for the identifier of the last entry. Redis’s $ (“entries appended after this call began”) would have avoided it, but there is a gap between two iterations of the loop, and everything written in that gap would land in neither — work would vanish without a trace. By resolving the position once, the handle turns the tail into an unbroken chain of identifiers.

History is deliberately not read through follow(): range() and after() are there for that.

`follow()` does not take a connection from the pool

A blocking read occupies a connection for the whole wait, so the handle opens its own. close() gives it back. For the same reason the handle is worth keeping for the whole loop rather than taking a fresh one each iteration.

Observing

info()

Returns the server’s summary of the log.

Syntax

php
public function info(): array

Parameters

None.

Returns

An array of whatever the server reports: length, first-entry, last-entry, last-generated-id, the number of groups and more. An empty array for a key that does not exist.

Example

php
$events->info()['length'];              // 1204
$events->info()['last-generated-id'];   // '1755600000123-0'

groups()

Returns every consumer group on this stream.

Syntax

php
public function groups(): array

Parameters

None.

Returns

An array of groups as the server sees them: name, consumers, pending, last-delivered-id, lag.

Example

php
foreach ($events->groups() as $group) {
  if ($group['pending'] > 1000) {
      $this->alert("group {$group['name']} is falling behind");
  }
}

lag is how many entries the group has not seen yet; together with pending those are two different kinds of lateness: the first is “not read”, the second “read but not acknowledged”.

The key as a whole

Entries have no lifetimes of their own — growth is bounded by trimming. A lifetime belongs to the key.

expireKey()

Gives the whole log a lifetime.

Syntax

php
public function expireKey(int $seconds): bool

Parameters

$seconds — in how many seconds to delete the key entirely, counted from now.

Returns

true if the lifetime was set; false if there is no key.

keyTtl()

Reports how long the log has left to live.

Syntax

php
public function keyTtl(): ?int

Parameters

None.

Returns

The remaining seconds, or null — both when there is no lifetime and when there is no key.

deleteKey()

Deletes the log along with its groups.

Syntax

php
public function deleteKey(): bool

Parameters

None.

Returns

true if the key existed.

It takes both the entries and the groups with their pending bookkeeping.

name()

Returns the key’s name with the prefix — the one the server sees.

Syntax

php
public function name(): string

Parameters

None.

Returns

The full key name. Needed for commands the handle does not wrap: those run through raw().

close()

Closes the connection follow() opened.

Always safe to call: if there was no blocking read the method does nothing.

Syntax

php
public function close(): void

Parameters

None.

Returns

Nothing.


Consumer groups

Without a group every reader sees all entries and tracks its own position. A group changes both halves: an entry goes to one consumer of the group, and the server holds it in that consumer’s pending list until it is acknowledged.

Hence the guarantee this is taken for: a worker took an entry and died — it is not lost. It is charged to that worker, it shows up in pending(), and another worker can take it over with claimStale().

php
$events->ensureGroup('mailers');                          // once, at startup
$group = $events->group('mailers', consumer: 'worker-1');

Creating

ensureGroup()

Creates a consumer group if it does not exist yet.

The stream is created too if needed, so a group can be declared at application startup, before the first entry.

Syntax

php
public function ensureGroup(string $name, string $from = '0'): bool

Parameters

$name — the group’s name.

$from — where the group starts reading. '0' by default — from the start of the log, so the group receives all the accumulated history. '$' — only entries appended after the group was created.

Returns

true if the group was created; false if it already existed — the latter is not an error, and calling it on every startup is safe.

Example

php
$events->ensureGroup('mailers');          // from the start of the log
$events->ensureGroup('metrics', '$');     // only what is new

Why a group is not created automatically by `consume()`

Then a typo in the name would become a new empty group that silently receives nothing. Requiring an explicit creation leaves a typo one outcome only — a NOGROUP error that names it.

group()

Returns a handle on a group on behalf of a particular consumer.

It creates nothing and does not go to the server — like the other handles, it is a view.

Syntax

php
public function group(string $name, string $consumer = 'default'): RedisStreamGroup

Parameters

$name — the group’s name. It must exist; ensureGroup() creates it.

$consumer — the name of the consumer the commands will act as. 'default' by default.

Returns

A RedisStreamGroup.

Example

php
$group = $events->group('mailers', consumer: 'worker-1');

RedisStreamGroup reference

RedisStreamGroup is a handle on a “group + consumer” pair. The consumer name in it is not incidental: the server keeps the pending list per consumer, not per group, so every recovery method acts on somebody’s behalf.

Consuming

consume()

Takes entries that have not been delivered to anyone in the group yet.

Every entry goes to exactly one consumer of the group and lands in that consumer’s pending list until it is acknowledged.

Syntax

php
public function consume(int $count = 10, float $timeout = 0.0): array

Parameters

$count — no more than this many entries per call. 10 by default.

$timeout — how many seconds to wait for new entries. 0.0 by default — wait indefinitely. A negative value means “do not wait at all”: look and return.

Returns

An array of StreamEntry; empty if there is nothing new.

Errors

RedisCommandException reading NOGROUP — if the group does not exist.

Example

php
$group->consume(count: 10, timeout: 5);   // waits, on its own connection
$group->consume(timeout: -1);             // look and return, connection from the pool

A negative timeout is the only mode that does not take a separate connection; it is for checks and one-off looks.

backlog()

Returns what this consumer has already taken but not acknowledged.

This is the first thing a restarted worker should do: entries taken before the restart are charged to it and will not arrive through consume() again. Reading them back is how it resumes without waiting for somebody to claim them on idleness.

Syntax

php
public function backlog(int $count = 10): array

Parameters

$count — no more than this many entries. 10 by default.

Returns

An array of StreamEntry. It does not wait: with nothing pending it returns an empty array immediately.

Example

main/Daemons/MailWorker.php
$group = $events->group('mailers', consumer: $this->workerName());

// 1. finish your own
foreach ($group->backlog() as $entry) {
  $this->handle($entry);
  $group->ack($entry);
}

// 2. and only then take new work
while (!$this->stopping) {
  foreach ($group->consume(count: 10, timeout: 5) as $entry) {
      $this->handle($entry);
      $group->ack($entry);
  }
}

$group->close();

ack()

Marks entries as handled, removing them from the pending list.

Syntax

php
public function ack(StreamEntry|string ...$entries): int

Parameters

...$entries — entries or their identifiers, mixed freely and as many as you like. You may pass none.

Returns

The number of entries that were pending and are now acknowledged. With no arguments — 0.

Example

php
$group->ack($entry);              // the object
$group->ack($entry->id);          // or the identifier
$group->ack(...$entries);         // or a batch

Acknowledging does not delete the entry from the log

ack() closes the books in the group: the entry stops being pending. The entry itself stays in the stream and is still visible in range() — it is a log, after all. Growth has to be bounded separately, by trimming.

The package will not acknowledge automatically on delivery: that would destroy the group’s one guarantee — that an entry outlives the worker that took it and did not finish.

Recovery

pending()

Returns entries the group delivered that nobody has acknowledged.

Syntax

php
public function pending(int $count = 100, ?string $consumer = null): array

Parameters

$count — no more than this many entries. 100 by default.

$consumer — restrict to one consumer. null by default — across the whole group.

Returns

An array of PendingEntry — not the entries themselves but their record cards: the identifier, who it is charged to, how long it has been idle and how many times it was delivered.

Example

php
foreach ($group->pending() as $stuck) {
  $stuck->id;           // '1755600000123-0'
  $stuck->consumer;     // 'worker-1'
  $stuck->idleMs;       // 61240 — how many milliseconds ago it was delivered
  $stuck->deliveries;   // 3 — how many times it was delivered
}

The two numbers answer different questions. A large idleMs with a single delivery means the consumer died holding the entry. A growing deliveries means the entry is taking down everyone who picks it up.

What to do about the second — restart, raise an alert, park it as dead — is the application’s decision: only the job’s author knows the cost of a retry.

pendingCount()

Reports how many entries in total the group delivered without receiving an acknowledgement.

Syntax

php
public function pendingCount(): int

Parameters

None.

Returns

The number of pending entries across the whole group.

claimStale()

Claims entries that have been sitting with another consumer for too long.

Syntax

php
public function claimStale(int $idle, int $count = 10, string $from = '0-0'): array

Parameters

$idle — how many milliseconds an entry must have been idle to be claimable. Pick the threshold from the longest honest processing time: if sending an email takes up to thirty seconds, idle: 60_000 will not take work away from a living worker.

$count — no more than this many entries. 10 by default.

$from — where in the pending list to start scanning. '0-0' by default — from the beginning.

Returns

An array of StreamEntry — entries already claimed and ready to handle.

Example

php
foreach ($group->claimStale(idle: 60_000) as $entry) {
  $this->handle($entry);
  $group->ack($entry);
}

A second call will not return the same entries

Claiming resets the idle counter, so the next call with the same threshold no longer sees them — the loop converges rather than spinning in place. The delivery counter does go up, and it shows that an entry has come round a second time.

consumers()

Returns every consumer of the group as the server sees them.

Syntax

php
public function consumers(): array

Parameters

None.

Returns

An array of consumers with the fields name, pending, idle.

Example

php
foreach ($group->consumers() as $consumer) {
  if ($consumer['idle'] > 300_000 && $consumer['pending'] > 0) {
      $this->alert("{$consumer['name']} has been silent for five minutes and is holding work");
  }
}

destroy()

Deletes the group along with its bookkeeping and consumers.

The log’s entries stay — only the record of who read what is removed.

Syntax

php
public function destroy(): bool

Parameters

None.

Returns

true if the group existed.

Housekeeping

as()

Returns the same group on behalf of a different consumer.

It does not change the original handle — it creates a new one. Needed where one process acts for several: walking somebody else’s stuck entries, for instance.

Syntax

php
public function as(string $consumer): RedisStreamGroup

Parameters

$consumer — the consumer’s name.

Returns

A new RedisStreamGroup with the same stream and group.

Example

php
$group->as('worker-2')->backlog();    // what the second worker did not finish

RedisStreamGroup::name()

Returns the group’s name.

Syntax

php
public function name(): string

Returns

The group’s name — what was passed to group().

consumer()

Returns the name of the consumer the handle acts as.

Syntax

php
public function consumer(): string

Returns

The consumer’s name.

Consumer names have to be unique

The server keeps the pending list per consumer. Two workers with the same name share one list and will pick up each other’s work — including what the first is processing right now. Name them after something stable and distinguishable: a worker slot, a host and a pid.

RedisStreamGroup::close()

Closes the connection opened by a blocking consume().

Syntax

php
public function close(): void

Returns

Nothing.

Mapping to Redis commands

Method Command
add() XADD, with a cap XADD ... MAXLEN
trim() / trimBefore() XTRIM MAXLEN / XTRIM MINID
delete() / count() XDEL / XLEN
range() / reverse() XRANGE / XREVRANGE
after() / follow() XREAD / XREAD BLOCK
info() / groups() / consumers() XINFO STREAM / GROUPS / CONSUMERS
ensureGroup() / destroy() XGROUP CREATE / XGROUP DESTROY
consume() / backlog() XREADGROUP with > / with 0
ack() XACK
pending() / pendingCount() XPENDING
claimStale() XAUTOCLAIM

What is not here — per-entry XCLAIM, XSETID, XGROUP CREATECONSUMER — is reachable through raw() together with name().

Next

  • Lists — a queue, when history is not needed
  • Hashes — an object’s state alongside the event log
  • Stores — the prefix, values, raw()
  • Connection pool — why blocking reads stay out of the pool