October 7, 2026 · 18 min read
Redis Pub/Sub vs Streams for AI Pipelines: Broadcast the News, Queue the Work
A plain-English guide to the two ways Redis moves messages between services. When a help page changes in a RAG app, the re-embedding work goes through a Stream and the cache invalidation goes out on pub/sub. Real crashes, lost messages, duplicate work and timeouts, with every output shown.
Redis has two ways to send a message from one service to another, and they look alike in code. Pub/sub announces something to everyone listening right now. Streams hand out work so that a task is kept until someone confirms it's done. Pick the wrong one and you get work done three times, or news nobody heard. This guide runs both on a real RAG pipeline, crashes workers on purpose, and shows what each one does.
You don't need to have used either. The caching guide is the background: it ended with a question this guide answers. Every section ends with a one-line Remember, and there's a summary at the end.
The guide has three parts:
- Part 1 · Understand it: what pub/sub and Streams are, and the one question that tells them apart, sections 1–5.
- Part 2 · Make the work reliable: timeouts, crashed workers, duplicate work and poison tasks, sections 6–10.
- Part 3 · Put it together: the full pipeline, how to choose, and Azure Managed Redis, sections 11–14.
Part 1
Understand it
Two kinds of message, and the Redis feature built for each.
1. One edit, two jobs
The running example is Parcelo again, the invented shop with a 12-page help centre and a RAG assistant. A support agent changes the refund policy on the returns page from 14 days to 30 days, and saves. Two things now have to happen.
Work. The page has to be cut into chunks and re-embedded, so the assistant can find the new text. That takes time and should happen exactly once: embedding the same page three times wastes money and can leave duplicate chunks.
News. Every running copy of the assistant has to forget its cached answers about refunds. If there are three copies, all three must hear it.
The editor shouldn't do either job itself. Saving a page should take milliseconds, not wait for an embedding model. And the editor shouldn't need to know how many copies of the assistant are running today. So the editor drops a message into Redis and moves on. The question is what kind of message.
2. Pub/sub: a radio station
Pub/sub (publish/subscribe) works like a radio station. The station broadcasts on a frequency. Every radio tuned in at that moment hears it. The station doesn't know who's listening, and a radio that's switched off misses the song.
In Redis the frequency is a channel, just a name like cache:invalidate. A subscriber tells Redis which channels it wants, then waits. A publisher sends a message to a channel, and Redis pushes a copy to every subscriber connected at that moment:
# every API instance, at startup
ps = r.pubsub(ignore_subscribe_messages=True)
ps.subscribe("cache:invalidate")
for message in ps.listen(): # waits; yields each message as it arrives
forget_answers_about(message["data"])
# the publisher, whenever something changes
receivers = r.publish("cache:invalidate", "returns") # returns how many heard it
PUBLISH returns how many subscribers received the message: a count, not a confirmation that anyone acted on it. You can also subscribe to a pattern: ps.psubscribe("cache:*") hears every channel that starts with cache:.
One rule surprises people: a connection that has subscribed can only receive messages. It can't run GET or SET any more, which is why redis-py makes you create a separate pubsub() object.
3. What pub/sub doesn't do
"Remembers nothing" sounds abstract until a subscriber restarts. In the lab, one API instance (api-2) was subscribed to cache:invalidate. It disconnected for a restart while two pages changed, then came back. I sent each change both ways: on the channel, and into a stream (section 4):
publish 'returns' -> 1 subscriber(s)
api-2 disconnects for a restart
publish 'shipping' -> 0 subscriber(s) (nobody listening)
publish 'lost' -> 0 subscriber(s)
api-2 is back and subscribed again
publish 'errors' -> 1 subscriber(s)
pub/sub: api-2 hears after restart -> ['errors']
stream: api-2 reads its group -> ['returns', 'shipping', 'lost', 'errors']
The publisher got 0 back and no error. Nothing anywhere recorded that shipping and lost were missed. api-2 kept serving its cached answers for those pages.
There's a second limit: a subscriber that can't keep up. Redis buffers messages for each subscriber, but only up to a point. I subscribed a client and then never read from it, and published 64 KB events:
client-output-buffer-limit: … pubsub 33554432 8388608 60
published 592 x 64 KB; last PUBLISH reached 0 subscriber(s)
pub/sub clients still connected: 0
The default limit for pub/sub clients is 32 MB (or 8 MB for 60 seconds). Past it, Redis closes the slow subscriber's connection. That protects Redis, and it means a slow consumer loses messages too.
4. Streams: a log with owners
A Redis Stream is an append-only log. New entries go on the end, each with an ID, and they stay until you trim them. Think of a board of job tickets: tickets are pinned up in order, a worker takes one and writes their name on it, and crosses it out when it's done.
r.xadd("docs:reindex", {"doc": "returns"}) # pin a ticket; returns its ID
On its own, anyone can read a stream from any point. The part that makes it a work queue is a consumer group: a named team of workers sharing one stream. Redis keeps the team's bookkeeping:
Stream: an append-only log, oldest on the left
Consumer group “embedders” (Redis keeps this for you)
- last-delivered-id = 1791379730278-0. Reading with
>means “give me entries after this”. - pending list: delivered but not acknowledged, with owner, idle time and delivery count.
- consumers: worker-a, worker-b. Created the first time each one reads.
An ID is the time in milliseconds when the entry was added, then a sequence number for entries added in the same millisecond.
- Each entry is delivered to one member of the group.
- A delivered entry goes on the group's pending list, with its owner, until that owner acknowledges it with
XACK. - Reading with the ID
>means "give me entries nobody in my group has had". Reading with0means "give me my own pending entries", which is how a restarted worker finds its unfinished work.
Because entries stay, the restarted api-2 in section 3 could read everything it missed. Its group remembered where it had got to.
5. The question that decides
Same three workers, same six "re-embed this page" jobs, sent both ways:
all 6 jobs sent
A) pub/sub: 3 workers subscribed to 'jobs', 6 jobs published
worker-1 received 6: ['returns', 'shipping', 'tracking', 'lost', 'address', 'errors']
worker-2 received 6: ['returns', 'shipping', 'tracking', 'lost', 'address', 'errors']
worker-3 received 6: ['returns', 'shipping', 'tracking', 'lost', 'address', 'errors']
=> 18 deliveries for 6 jobs
B) Stream + consumer group 'embedders': same 3 workers, same 6 jobs
worker-1 processed 2: ['returns', 'lost']
worker-2 processed 2: ['shipping', 'address']
worker-3 processed 2: ['tracking', 'errors']
=> 6 deliveries for 6 jobs
With pub/sub, adding workers adds copies, not capacity: every page gets embedded three times. With a consumer group, adding a worker adds capacity, and you write no load-balancing code.
Neither is better. They answer different questions:
Every receiver: it's news, use pub/sub. One of them: it's work, use a stream with a consumer group.
Part 2
Make the work reliable
What goes wrong with workers, and the Streams features that handle it.
6. The worker loop
Every embedder runs the same loop. Here it is, with the three rules that matter marked:
try:
r.xgroup_create("docs:reindex", "embedders", id="0", mkstream=True)
except redis.ResponseError as e:
if "BUSYGROUP" not in str(e): # rule 1: every copy runs this; only the first creates it
raise
me = f"{socket.gethostname()}-{os.getpid()}"
while True:
reply = r.xreadgroup("embedders", me, {"docs:reindex": ">"}, count=1, block=5000)
for _stream, entries in reply: # one list per stream, then the entries
for task_id, fields in entries:
reembed(fields["doc"])
r.xack("docs:reindex", "embedders", task_id) # rule 2: only after the work succeeded
block=5000 tells Redis: if there's nothing new, wait up to 5 seconds before answering. That keeps idle workers from spinning. It also hides rule 3, which I found by running the loop with redis-py's defaults:
redis-py 8.0.1, default socket_timeout = 5
defaults block=5000 on an empty stream -> TimeoutError: Timeout reading from socket after 58.93 s
socket_timeout=30 block=5000 on an empty stream -> [] after 5.08 s
The client waits at most 5 seconds for a reply (socket_timeout). Redis was told to wait 5 seconds before replying. The client gives up first, retries (redis-py 8 retries failed reads by default), gives up again, and raises after almost a minute. Rule 3: the client's socket timeout must be longer than any block you ask for. With socket_timeout=30, the empty read came back after 5.08 seconds as an empty list, and the loop simply goes round again.
7. When a worker dies
This is the reason to use Streams for work. I started a real worker process that took one task and then killed itself with os._exit(137), the same exit code as an out-of-memory kill. No cleanup ran:
6. worker-b acknowledges. The pending list is empty. Nothing was lost.
[worker-a pid 2854] took returns (1791379730277-0) ... killed before XACK
1) worker-a process exited with code 137
2) XPENDING: 1 pending, held by [{'name': 'worker-a', 'pending': 1}]
1791379730277-0 owner=worker-a idle=1107 ms delivered=1x
3) worker-b XAUTOCLAIMed returns (1791379730277-0) idle > 1000 ms, re-embedded it, XACK
4) worker-b took the next new task shipping, XACK
5) pending left: 0
The task survived the crash. XPENDING shows what's stuck, who holds it and for how long. XAUTOCLAIM finds tasks idle longer than a limit and moves them to a healthy worker in one command:
_next, claimed, _deleted = r.xautoclaim("docs:reindex", "embedders", "worker-b",
min_idle_time=60_000, start_id="0-0")
I used 1 second so the demo didn't wait. In production, set it longer than your slowest normal task, or a healthy but slow worker will have its task taken away mid-job.
There's a second recovery path. If the crashed worker restarts with the same consumer name, it can read its own pending tasks with ID 0 before taking new ones. That covers the common crash-and-restart. XAUTOCLAIM covers the worker that never comes back.
With pub/sub there is no equivalent. A subscriber that dies mid-task leaves no trace.
8. At least once, not exactly once
Section 7's worker died before doing the work. What if it dies after doing the work but before XACK? Redis can't tell the difference. The task is still pending, so it gets claimed and done again.
I ran exactly that. The work is "write the page's chunks". The worker wrote them, died, and the task was reclaimed:
append (naive) chunks for 'returns' after crash + reclaim: 9 (should be 3)
replace (idempotent) chunks for 'returns' after crash + reclaim: 3 (should be 3)
append (naive)
9 chunks (should be 3)
The old chunks were never removed, and the reclaimed task appended the new ones a second time.
replace (idempotent)
3 chunks (should be 3)
Delete the page's chunks, then write them. Running it twice gives the same result as running it once.
The naive version appended chunks, so the old 14-day chunks stayed and the new 30-day ones went in twice. The assistant would now retrieve both policies and could quote either one.
So a consumer group gives you at-least-once delivery: every task is done, sometimes more than once. The fix is on your side. Make the work idempotent, meaning that running it twice leaves the same result as running it once. For re-embedding, that means "replace this page's chunks", not "add chunks". For other work, use upserts keyed by the task or document ID, or record finished task IDs and skip them.
9. Poison tasks
Some tasks fail every time: a corrupt PDF, a page the parser can't read. Left alone, they get claimed, fail, wait, and get claimed again forever. The pending list tracks how many times each task was delivered, so the worker can count:
round 1: giftcards failed on delivery 1, left pending for a retry
round 1: export re-embedded on delivery 1, XACK
round 2: giftcards failed on delivery 2, left pending for a retry
round 3: giftcards failed on delivery 3 -> moved to docs:dead, XACK
pending: 0 dead letters: [('1791379816171-0', {'doc': 'giftcards', 'path': '/broken.pdf', 'orig_id': '1791379816168-0', 'error': 'PDF parse failed', 'deliveries': '3'})]
After 3 failed deliveries, the worker copies the task into a separate dead-letter stream, docs:dead, with the error, and acknowledges the original. The queue keeps moving, and a person can look at docs:dead later. XPENDING with a range (xpending_range in redis-py) gives you the count in its times_delivered field.
10. Keep the stream bounded
XACK removes a task from the pending list. It does not delete the entry from the stream. A stream that gets an entry for every page edit grows forever unless you trim it:
r.xadd("docs:reindex", {"doc": "returns"}, maxlen=10_000, approximate=True)
maxlen keeps roughly the newest 10,000 entries. approximate=True lets Redis trim in efficient chunks, so the length can sit a little above the limit. Choose the limit so that unprocessed entries are never trimmed: it should be far larger than your worst backlog. The entries you keep are also a free history of what was re-indexed and when.
Part 3
Put it together
The full pipeline from one edit to a fresh answer, then how to choose.
11. The whole pipeline
Here's the edit from section 1, end to end, on the real app from the caching guide: real pgvector, real embeddings, real cache keys. The answer cache lives in two places, like many production apps: a shared Redis cache, and a small in-memory cache inside each of three API instances.
7. The next question misses every cache and gets the new policy.
1) before the edit, cached answer: Parcelo Plus membership fees are refundable only within 14 days of pur…
2) CMS saves the page and queues work: XADD docs:reindex -> 1791379825171-0
embedder-1 deleted 2 shared Redis key(s) built from 'returns'
3) embedder-1 re-embedded 'returns' in 63 ms, XACK, PUBLISH cache:invalidate -> 3 API instances
4) api-1 heard 'returns' 66 ms after the edit, dropped 1 entry from its in-process cache
4) api-2 heard 'returns' 66 ms after the edit, dropped 1 entry from its in-process cache
4) api-3 heard 'returns' 66 ms after the edit, dropped 1 entry from its in-process cache
5) next question: answer miss, retrieval miss -> Parcelo Plus membership fees are refundable only within 30 days of pur…
Three details are worth copying:
- The order. The embedder re-embeds first, then deletes cache keys, then announces. Delete before re-embedding, and a request in between rebuilds the cache from the old chunks.
- One delete for the shared cache. Every instance reads the same Redis, so the worker deletes those keys once (using the dependency sets from the caching guide). No broadcast is needed for that.
- A broadcast for the local caches. Each instance's in-memory cache is invisible to the others. Only pub/sub reaches all of them, and it did in 66 ms.
The broadcast has pub/sub's weakness from section 3: an instance that is restarting misses it. Here that's acceptable, because a restarting instance starts with an empty in-memory cache anyway, and every local entry also has a short TTL. If missing an announcement were not acceptable, the next section says what to use.
12. How to choose
| You need… | Use | Because |
|---|---|---|
| Every running copy to react (cache invalidation, config reload, live status) | Pub/sub | One publish reaches all current subscribers, sub-millisecond, no bookkeeping |
| One worker to do each task (embedding, OCR, model calls) | Stream + one consumer group | Each entry goes to one worker; scale by adding workers |
| No task lost if a worker crashes | Stream + consumer group | Pending list, XAUTOCLAIM, delivery counts |
| Several different services to each see every event, reliably | Stream with one consumer group per service | Each group has its own position and pending list, so a restarting service catches up |
| A live feed to browsers | Pub/sub behind a WebSocket server | The browser can't talk to Redis; the server subscribes and forwards |
The fourth row is the one people miss. If the analytics service and the billing service must both see every "document processed" event, and neither may miss one, don't use pub/sub. Give each its own consumer group on the same stream.
Two things to weigh with Streams: they're a little more code (groups, acknowledgements, claiming), and you manage their memory (section 10). Two things to weigh with pub/sub: nothing is stored, and nothing confirms that a receiver acted.
13. On Azure Managed Redis
Pub/sub and Streams are standard Redis features, so the commands and the code above work unchanged on Azure Managed Redis. Only the connection changes, as in the caching guide: port 10000, TLS, and a Microsoft Entra ID token from the redis-entraid package.
Two settings matter more for messaging than for caching:
socket_timeoutmust be longer than yourblocktime, on Azure as locally (section 6).- Long-lived connections. A subscriber or a blocked worker holds its connection open for hours. The Entra ID credential provider refreshes the token in the background, so the connection stays authorised.
I didn't run this lab on Azure for this article. The connection details follow the Microsoft Learn module "Implement event messaging with Azure Managed Redis" (AI-200), checked in October 2026.
14. Common mistakes
- Pub/sub for work. Section 5: every worker does every job. Adding workers adds duplicates.
- Pub/sub for things that must not be missed. Section 3: a restart or a slow consumer drops messages silently, and
PUBLISHreturns a count, not a guarantee. XACKbefore the work. If the worker crashes after acknowledging, the task is gone. Acknowledge last.- Non-idempotent tasks. Section 8: a reclaimed task ran twice and left 9 chunks instead of 3.
- No claim loop. Without
XAUTOCLAIM(or a restart with the same consumer name), a dead worker's tasks sit in the pending list forever. - A claim time shorter than the work. A slow but healthy worker loses its task to another worker, and the work runs twice.
socket_timeout≤block. Section 6: TimeoutError after almost a minute instead of an empty reply after 5 seconds.- Random consumer names on restart. A worker that restarts as a new name can't find its own pending tasks with ID
0. Use a stable name per instance. - Unbounded streams.
XACKdoesn't delete. UseMAXLEN. - Reading the reply with one loop.
xreadgroupreturns a list per stream, then the entries. Loop over both.
The whole guide in ten lines
- A change creates work (one worker, once) and news (every running copy).
- Pub/sub delivers each message to every subscriber connected at that moment, and stores nothing.
- A restarting or slow subscriber misses messages, and nobody is told.
- A stream keeps entries; a consumer group gives each to one worker and tracks it until XACK.
- Same 6 jobs, 3 workers: pub/sub made 18 deliveries, a consumer group made 6.
- Set the socket timeout longer than the block time, or a waiting worker fails after almost a minute.
- A dead worker's task waits in the pending list; XAUTOCLAIM gives it to a live worker.
- Delivery is at least once: a reclaimed task ran twice, so make tasks idempotent.
- Count deliveries and move poison tasks to a dead-letter stream; cap streams with MAXLEN.
- The pipeline: queue the work in a stream, then announce the result on a channel, in that order.
Checked against Redis 7.4.11 (Docker image redis:7.4) and redis-py 8.0.1, with pgvector 0.8.7 on PostgreSQL 17.11 and all-MiniLM-L6-v2 for the pipeline in section 11, as of October 2026. Every output block is real output from my laptop on 7 October 2026; stream IDs and timings will differ on yours. The crash in sections 7 and 8 is a real process killed with os._exit(137). The "work" in section 8 writes list entries as a stand-in for chunk rows; the re-embedding in section 11 is real. The Azure section was not run and follows the Microsoft Learn AI-200 module. The Parcelo help centre is invented.