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.
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.
$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.
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:
$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
public function add(
array $fields,
?int $cap = null,
bool $exact = false,
string $id = '*',
): stringParameters
$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:
RedisCommandException: ERR The ID specified in XADD is equal or smaller
than the target stream top itemExample
$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
public function trim(int $cap, bool $exact = false): intParameters
$cap — how many entries to keep.
$exact — trim exactly rather than approximately. false by default.
Returns
The number of entries deleted.
Example
$events->trim(10_000); // cheap, the length is "about that"
$events->trim(10_000, exact: true); // exactly that many, more expensiveAn 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
public function trimBefore(string $id, bool $exact = false): intParameters
$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
$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
public function delete(string ...$ids): intParameters
...$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
public function count(): intParameters
None.
Returns
The number of entries; 0 for a key that does not exist.
range()
Returns entries in append order, taking nothing away.
Syntax
public function range(?string $from = null, ?string $to = null, ?int $count = null): arrayParameters
$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
$events->range(); // everything
$events->range(count: 100); // the first hundred
$events->range($since, count: 50); // from a known positionreverse()
Returns entries newest first.
A cheap way to see what just happened without reading the log from the start.
Syntax
public function reverse(?int $count = null): arrayParameters
$count — no more than this many entries. null by default — the whole log.
Returns
An array of StreamEntry, freshest first.
Example
$events->reverse(count: 10); // the last ten, newest on topafter()
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
public function after(string $id, int $count = 10): arrayParameters
$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
$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
public function follow(float $timeout = 0.0, int $count = 10, ?string $from = null): arrayParameters
$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
$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
public function info(): arrayParameters
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
$events->info()['length']; // 1204
$events->info()['last-generated-id']; // '1755600000123-0'groups()
Returns every consumer group on this stream.
Syntax
public function groups(): arrayParameters
None.
Returns
An array of groups as the server sees them: name, consumers, pending,
last-delivered-id, lag.
Example
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
public function expireKey(int $seconds): boolParameters
$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
public function keyTtl(): ?intParameters
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
public function deleteKey(): boolParameters
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
public function name(): stringParameters
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
public function close(): voidParameters
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().
$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
public function ensureGroup(string $name, string $from = '0'): boolParameters
$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
$events->ensureGroup('mailers'); // from the start of the log
$events->ensureGroup('metrics', '$'); // only what is newWhy 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
public function group(string $name, string $consumer = 'default'): RedisStreamGroupParameters
$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
$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
public function consume(int $count = 10, float $timeout = 0.0): arrayParameters
$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
$group->consume(count: 10, timeout: 5); // waits, on its own connection
$group->consume(timeout: -1); // look and return, connection from the poolA 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
public function backlog(int $count = 10): arrayParameters
$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
$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
public function ack(StreamEntry|string ...$entries): intParameters
...$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
$group->ack($entry); // the object
$group->ack($entry->id); // or the identifier
$group->ack(...$entries); // or a batchAcknowledging 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
public function pending(int $count = 100, ?string $consumer = null): arrayParameters
$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
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
public function pendingCount(): intParameters
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
public function claimStale(int $idle, int $count = 10, string $from = '0-0'): arrayParameters
$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
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
public function consumers(): arrayParameters
None.
Returns
An array of consumers with the fields name, pending, idle.
Example
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
public function destroy(): boolParameters
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
public function as(string $consumer): RedisStreamGroupParameters
$consumer — the consumer’s name.
Returns
A new RedisStreamGroup with the same stream and group.
Example
$group->as('worker-2')->backlog(); // what the second worker did not finishRedisStreamGroup::name()
Returns the group’s name.
Syntax
public function name(): stringReturns
The group’s name — what was passed to group().
consumer()
Returns the name of the consumer the handle acts as.
Syntax
public function consumer(): stringReturns
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
public function close(): voidReturns
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