Valkey has two ways to pass messages between processes, and they make opposite promises. Pub/sub is fire-and-forget: a message goes to whoever is subscribed at that instant and is then gone, with no storage, no acknowledgement and no replay. Streams are an append-only log: messages are stored, each has an id, consumers can read from any point, and consumer groups track which messages each worker has handled and which it has not acknowledged.
The rule that falls out of that is short. Use pub/sub when a missed message does not matter because the next one replaces it - live notifications, cache invalidation hints, presence. Use streams when every message must be processed at least once - orders, events that update other systems, anything you would be upset to lose. Most bugs with either come from using pub/sub where a stream was needed. This post covers both properly: the commands, the guarantees, consumer groups, the pending list, trimming, and the failure cases.
The difference in one table#
| Pub/sub | Streams | |
|---|---|---|
| Stored | No | Yes, until trimmed |
| Delivery | At most once, to current subscribers | At least once with consumer groups |
| Late subscriber sees old messages | No | Yes, from any id |
| Acknowledgement | None | XACK per message |
| Load balancing across workers | No, every subscriber gets every message | Yes, within a consumer group |
| Survives a Valkey restart | Nothing to survive | Yes, with persistence |
| Memory | Only buffers for slow subscribers | Grows with retained entries |
| Typical use | Live updates, invalidation, chat fan-out | Event logs, work distribution, audit trails |
There is a third option people forget: a plain list with LPUSH and BRPOP, which is what simple job queues use. It delivers each message to one consumer and stores it until taken, but once a consumer has popped a message, a crash loses it. Streams with consumer groups fix exactly that, which is why they exist. For full job-queue libraries built on these primitives, see job queues with Valkey.
Pub/sub: commands and behaviour#
A publisher sends to a channel; every connection currently subscribed to that channel receives it.
# connection ASUBSCRIBE orders:updatesPSUBSCRIBE chat:room:*# connection BPUBLISH orders:updates '{"id":512,"status":"shipped"}'(integer) 1PUBLISH returns the number of subscribers that received the message. A return of 0 means the message went nowhere - nobody was listening, and it is gone. PSUBSCRIBE subscribes to a glob pattern; useful, but every published message is matched against every pattern subscription, so thousands of patterns cost CPU on every publish.
Things that surprise people:
- A subscribed connection is special. Under the RESP2 protocol, a connection in subscribe mode can only issue subscribe-related commands and
PING. Client libraries therefore use a separate connection for subscribing. Do not try to share your main client. - Disconnected means missed. A subscriber that reconnects after a network blip, a deploy or a garbage-collection pause receives nothing that was published while it was away. There is no catch-up.
- Slow subscribers are disconnected. Messages waiting for a subscriber are buffered in Valkey's memory. The
client-output-buffer-limitfor pub/sub clients defaults to32mb 8mb 60: a subscriber whose buffer exceeds 32 MB, or stays above 8 MB for 60 seconds, is disconnected. That protects the server from one stuck consumer, and it means a slow subscriber loses messages. - Pub/sub is not stored, so persistence does nothing for it. AOF and snapshots save keys; channels are not keys.
Sharded pub/sub (SSUBSCRIBE, SPUBLISH) exists for cluster deployments, where it keeps messages within one shard. On a single server it behaves like ordinary pub/sub.
What pub/sub is good for#
Given those properties, pub/sub fits work where losing a message is harmless:
- Real-time fan-out to browsers. Several web processes each hold WebSocket connections; an event published once reaches every process, which forwards it to its connected users. Socket.IO's Redis adapter and ASP.NET Core SignalR's Redis backplane both work this way - SignalR real-time apps and WebSockets behind a reverse proxy cover those setups. If a user misses a "new message" ping, they see the message on their next refresh.
- Cache invalidation hints. Processes with in-memory caches subscribe to
cache:invalidate; a write publishes the key that changed. Combine with a TTL on the local cache so a missed hint only means briefly stale data. - Presence and typing indicators. By definition only relevant to whoever is listening right now.
- Development and debugging.
SUBSCRIBEinvalkey-clito watch events flow.
import { createClient } from "redis";const pub = createClient({ url: process.env.VALKEY_URL });const sub = pub.duplicate();await Promise.all([pub.connect(), sub.connect()]);await sub.subscribe("orders:updates", (message) => { broadcastToBrowsers(JSON.parse(message));});await pub.publish("orders:updates", JSON.stringify({ id: 512, status: "shipped" }));Note duplicate(): a second connection for the subscriber, for the reason above.
Keyspace notifications
Valkey can also publish events about its own keys: a key expiring, being deleted, being written. These are keyspace notifications, delivered over ordinary pub/sub channels such as __keyevent@0__:expired. They are switched off by default, because generating them costs CPU on every write, and they are enabled with the notify-keyspace-events setting - Ex for expiry events, for example.
They are tempting for "do something when this session expires" or "run this job when the timer key expires", and they inherit every weakness of pub/sub. A listener that is disconnected when the event fires never hears about it. Expiry events also fire when Valkey actually removes the key, which happens when the key is accessed or found by the background expiry cycle - not necessarily at the exact second the TTL ran out. For anything that must happen, a sorted set scored by due time, polled by a worker, is the reliable version of the same idea. On a hosted instance, changing notify-keyspace-events needs CONFIG SET permission, which your user may not have.
Streams: an append-only log#
A stream is a key holding an ordered sequence of entries. Each entry is a small set of field-value pairs with an id of the form <milliseconds>-<sequence>, assigned by the server.
XADD orders * type created id 512 total 49.90"1791468000123-0"XLEN orders(integer) 1XRANGE orders - +1) 1) "1791468000123-0" 2) 1) "type" 2) "created" 3) "id" 4) "512" 5) "total" 6) "49.90"The * asks Valkey to generate the id. XRANGE with - and + reads everything from the oldest to the newest; XREVRANGE reads backwards. Because ids are timestamps, you can read "everything since 09:00" by giving a millisecond value as the start.
Without consumer groups, any number of readers can follow a stream independently with XREAD, each remembering the last id it saw:
XREAD COUNT 100 BLOCK 5000 STREAMS orders 1791468000123-0BLOCK 5000 waits up to five seconds for new entries; $ instead of an id means "only entries added from now on". This is already a big step past pub/sub: a reader that disconnects can resume from its last id and miss nothing, as long as the entries have not been trimmed. Every reader sees every entry, so this is fan-out with replay.
A few design choices are worth making at this stage. Keep entries small and self-describing: an event type, the id of the thing that changed and perhaps a version, rather than the whole record - consumers can load the current state from the database, and small entries keep the stream cheap to retain. Use one stream per kind of event (orders, payments) rather than one stream for everything, so consumers that care about orders do not have to read and skip every payment. And let the server assign ids; supplying your own is possible but they must always increase, which is easy to get wrong across several producers.
Consumer groups: sharing work with acknowledgements#
To spread entries across several workers so each entry is handled once, create a consumer group. The group tracks a last-delivered id for the group as a whole, plus a pending entries list (PEL) of messages delivered to a consumer but not yet acknowledged.
XGROUP CREATE orders billing $ MKSTREAMXREADGROUP GROUP billing worker-1 COUNT 10 BLOCK 5000 STREAMS orders >XACK orders billing 1791468000123-0What each of those three lines does:
XGROUP CREATE ... $starts the group at the end of the stream; use0to process everything already in it.MKSTREAMcreates the stream if it does not exist yet.XREADGROUPwith the special id>asks for entries never delivered to anyone in the group. Each entry goes to exactly one consumer. The consumer name (worker-1) is arbitrary but should be stable across restarts of the same worker.XACKremoves the entry from the pending list. Until then, the group remembers thatworker-1has it.
Several groups can read the same stream independently: billing and email each get every order, and within each group the work is split among consumers. That is the pattern for an event log that several subsystems react to.
import os, redisr = redis.Redis.from_url(os.environ["VALKEY_URL"], decode_responses=True)GROUP, NAME = "billing", os.environ.get("WORKER_NAME", "worker-1")try: r.xgroup_create("orders", GROUP, id="0", mkstream=True)except redis.ResponseError as e: if "BUSYGROUP" not in str(e): raisewhile True: # first, anything delivered to me before a crash and never acknowledged for stream, entries in r.xreadgroup(GROUP, NAME, {"orders": "0"}, count=50) or []: for entry_id, fields in entries: handle(fields); r.xack("orders", GROUP, entry_id) for stream, entries in r.xreadgroup(GROUP, NAME, {"orders": ">"}, count=50, block=5000) or []: for entry_id, fields in entries: handle(fields); r.xack("orders", GROUP, entry_id)Reading with id 0 instead of > returns this consumer's own pending entries - what it received before it crashed and never acknowledged. Processing those first on startup is how a restarted worker finishes what it started.
Delivery guarantees and the pending list#
Streams with consumer groups give at-least-once delivery. A consumer that processes an entry and crashes before XACK will see it again; so will another consumer that claims it. Handlers must therefore be idempotent - safe to run twice - exactly as with any queue.
Entries whose consumer died permanently need to be claimed by someone else:
XPENDING orders billingXPENDING orders billing - + 10XAUTOCLAIM orders billing worker-2 60000 0-0 COUNT 10XPENDING summarises what is outstanding and, with a range, lists entries with their owner, idle time and delivery count. XAUTOCLAIM transfers entries idle for longer than 60,000 ms to worker-2 and returns them. Run it periodically in each worker, and you have automatic recovery from dead consumers.
The delivery count is your poison-message detector. An entry delivered ten times is not going to succeed on the eleventh; copy it to a dead-letter stream with XADD orders:dead ..., XACK it from the original group, and alert. Without that step one malformed message cycles through your workers forever.
Clean up consumers that no longer exist with XGROUP DELCONSUMER, after claiming their pending entries - deleting a consumer discards its pending list.
Trimming, memory and persistence#
A stream keeps every entry until you remove it. Acknowledging an entry does not delete it - XACK only updates the group's pending list. Without trimming, a busy stream grows until memory runs out.
XADD orders MAXLEN ~ 100000 * type created id 513XTRIM orders MAXLEN ~ 100000XTRIM orders MINID ~ 1790863200000MAXLEN caps the number of entries; MINID removes everything older than an id, which is a time-based retention ("keep seven days"). The ~ makes trimming approximate, letting Valkey remove whole internal nodes at once; it may keep slightly more than the cap and is much cheaper than exact trimming. Use it unless you need an exact count.
Trimming does not check whether every group has processed an entry. If a consumer group falls far enough behind, trimming deletes entries it never saw. Size retention generously relative to how long a consumer can be down, and monitor group lag - XINFO GROUPS orders shows each group's last-delivered id, pending count and, on current versions, a lag figure.
Memory per entry depends on field count and value sizes; small entries cost tens of bytes each because streams store them compactly. A hundred thousand small events typically fit in a few tens of megabytes - measure with MEMORY USAGE orders rather than trusting any estimate. Unlike pub/sub, streams are keys, so AOF and snapshots persist them, and the eviction policy applies to them. A stream that must not lose entries needs noeviction, exactly like a queue; Valkey persistence: AOF and RDB covers the loss window on a hard stop.
On RE:NODE a Valkey plan saves to disk with both AOF and snapshots, so stream entries and consumer-group state survive a restart. It is reached directly on its host and port with a password generated for that server; there is no proxy slot in front, which is what long-lived blocking reads need.
Troubleshooting#
Subscribers miss messages. Expected with pub/sub when they are disconnected, slow or not yet subscribed. If that is not acceptable, the design needs a stream.
`PUBLISH` returns 0. No subscribers on that channel at that moment. Check the channel name - pub/sub channels are case sensitive and have nothing to do with logical database numbers.
Consumer group lag keeps growing. Consumers are slower than producers. Add consumers to the group, raise COUNT, or find the slow handler.
Pending list grows without bound. Handlers are not calling XACK, or are crashing before it. Check XPENDING for owners and delivery counts.
`BUSYGROUP Consumer Group name already exists`. Harmless at startup; catch it, as in the example.
Memory climbs steadily with a stream. No trimming. Add MAXLEN ~ to XADD, or run XTRIM on a schedule.
FAQ#
Can I get pub/sub messages that were sent while my app was down?
No. Pub/sub stores nothing. If you need that, publish to a stream instead and have each reader track its last id; readers that were down resume where they left off.
Is a Valkey stream a replacement for Kafka?
For a single application or a handful of services, it covers the same pattern - an ordered log with independent consumer groups - with far less to operate. Kafka is built for retention measured in terabytes, partitions across many machines and replay over weeks. If your stream fits comfortably in memory, Valkey is the proportionate tool.
Should I use streams or a job queue library?
A job queue library when the messages are tasks with retries, delays and scheduling - BullMQ, Celery and RQ give you those already. Streams when the messages are events several consumers react to, or when you want an ordered history you can replay.
How do I make processing exactly-once?
You cannot with the transport alone; streams deliver at least once. Make the handler idempotent: record the entry id, or a business key, in your database inside the same transaction as the work, and skip anything already recorded.
Do pub/sub messages count towards memory limits?
Only while they are buffered for slow subscribers, and those buffers are capped per client by client-output-buffer-limit. Stream entries count fully, because they are stored data.




Comments
Completely anonymous: no account, no email, no cookie. We store the name you type, the text and the time - nothing else. Links are limited and markup is not rendered.