Aryan Tripathi — Writing
← All writing

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#streams#pub-sub#rag#ai-engineering#azure

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.

RememberA change usually creates two jobs: work that one worker must do once, and news that every running copy must hear.

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.

RememberPub/sub delivers a copy of each message to every subscriber connected at that moment, and remembers nothing.

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):

Animatedapi-2 restarts mid-stream · when I ran it
api-2 restartingPub/subStreamreturns✓heard✓readshipping✗lost✓readlost✗lost✓readerrors✓heard✓read
pub/sub: api-2 hears → ['returns', 'errors']stream: api-2 reads → ['returns', 'shipping', 'lost', 'errors']
step 6 of 6
Pub/sub only reaches subscribers that are connected right now. A stream keeps every entry, so a consumer group catches up after a restart.
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.

RememberPub/sub has no memory: a subscriber that is restarting, or too slow, misses messages without anyone being told.

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:

Diagramstream docs:reindex · group embedders

Stream: an append-only log, oldest on the left

1791379730277-0
doc: returns
pending · worker-a
1791379730278-0
doc: shipping
pending · worker-b
1791379730279-0
doc: lost
new
next XADD →

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.

The stream holds the entries. The consumer group remembers how far it has handed out work, and which entries were delivered but not yet acknowledged: the pending list.
  • 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 with 0 means "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.

RememberA stream keeps every entry; a consumer group gives each entry to one worker and tracks it until the worker acknowledges it.

5. The question that decides

Same three workers, same six "re-embed this page" jobs, sent both ways:

Animated3 workers · 6 re-embedding jobs · when I ran it
Pub/sub: broadcastchanneljobsworker-16 jobsworker-26 jobsworker-36 jobs18 deliveries for 6 jobsStream + consumer groupstreamjobs:queueworker-12 jobsworker-22 jobsworker-32 jobs6 deliveries for 6 jobs

all 6 jobs sent

step 7 of 7
Pub/sub copies every message to every subscriber: 18 deliveries, so every page would be re-embedded three times. A consumer group hands each entry to one worker: 6 deliveries, load shared.
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:

Should every receiver handle this message, or exactly one of them?

Every receiver: it's news, use pub/sub. One of them: it's work, use a stream with a consumer group.

RememberEvery receiver should act on it: pub/sub. Exactly one should: 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
Real outputXREADGROUP … BLOCK 5000 on an empty stream · when I ran it
0 s15 s30 s45 s60 sredis-py 8.0.1 defaultssocket_timeout = 5 s58.9 s → TimeoutError: Timeout reading from socketsocket_timeout = 30longer than the block5.1 s → [] (nothing new, loop again)
The worker asks Redis to wait up to 5 seconds for new work. With redis-py's default 5-second socket timeout, the client gives up at the same moment, then retries, and finally raises after almost a minute.

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.

RememberCreate the group idempotently, acknowledge only after the work succeeds, and set the socket timeout longer than the block time.

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:

Animateda real worker process killed mid-task · when I ran it
Streamworker-a ✗Pending listworker-bDonereturns…730277-0

6. worker-b acknowledges. The pending list is empty. Nothing was lost.

XACK docs:reindex embedders 1791379730277-0 → pending left: 0
step 6 of 6
A consumer group gives every delivered task an owner until it is acknowledged. If the owner dies, the task waits in the pending list until another worker claims it.
   [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.

RememberA task delivered to a dead worker waits in the pending list; XAUTOCLAIM hands it to a live one.

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)
Real outputwork done, then crash before XACK · when I ran it

append (naive)

9 chunks (should be 3)

14 days (old)Standard returnsGift cards30 daysStandard returnsGift cards30 daysStandard returnsGift cards

The old chunks were never removed, and the reclaimed task appended the new ones a second time.

replace (idempotent)

3 chunks (should be 3)

30 daysStandard returnsGift cards

Delete the page's chunks, then write them. Running it twice gives the same result as running it once.

A consumer group delivers each task at least once, not exactly once. A worker that finishes the work and dies before XACK will see its task done again. Make the work safe to repeat.

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.

RememberStreams deliver at least once. Make every task safe to repeat: replace, don't append.

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.

RememberCount deliveries, and after a few failures move the task to a dead-letter stream instead of retrying forever.

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.

RememberAcknowledging doesn't delete. Cap every stream with MAXLEN so memory stays bounded.

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.

Animatedthe whole path, real timings · when I ran it
CMSpage savedStreamdocs:reindexembedder-1group embeddersShared cacheRedis keysChannelcache:invalidateapi-1in-process cacheapi-2in-process cacheapi-3in-process cache

7. The next question misses every cache and gets the new policy.

answer miss, retrieval miss → “…within 30 days…”
step 7 of 7
The hybrid pattern: the work goes through a stream so exactly one worker owns it, and the news goes through pub/sub so every running copy hears it.
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.

RememberQueue the work in a stream, announce the result on a channel, and do them in that order.

12. How to choose

You need…UseBecause
Every running copy to react (cache invalidation, config reload, live status)Pub/subOne publish reaches all current subscribers, sub-millisecond, no bookkeeping
One worker to do each task (embedding, OCR, model calls)Stream + one consumer groupEach entry goes to one worker; scale by adding workers
No task lost if a worker crashesStream + consumer groupPending list, XAUTOCLAIM, delivery counts
Several different services to each see every event, reliablyStream with one consumer group per serviceEach group has its own position and pending list, so a restarting service catches up
A live feed to browsersPub/sub behind a WebSocket serverThe 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.

RememberPub/sub for "everyone, now, best effort". Streams for "someone, eventually, guaranteed", and one group per service when every service needs everything.

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_timeout must be longer than your block time, 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.

RememberOn Azure Managed Redis the messaging code is the same; set the socket timeout above the block time and use the Entra ID provider for long connections.

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 PUBLISH returns a count, not a guarantee.
  • XACK before 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. XACK doesn't delete. Use MAXLEN.
  • Reading the reply with one loop. xreadgroup returns a list per stream, then the entries. Loop over both.
RememberAcknowledge last, make tasks repeatable, claim stuck tasks, and never trust pub/sub with something that must happen.

The whole guide in ten lines

  1. A change creates work (one worker, once) and news (every running copy).
  2. Pub/sub delivers each message to every subscriber connected at that moment, and stores nothing.
  3. A restarting or slow subscriber misses messages, and nobody is told.
  4. A stream keeps entries; a consumer group gives each to one worker and tracks it until XACK.
  5. Same 6 jobs, 3 workers: pub/sub made 18 deliveries, a consumer group made 6.
  6. Set the socket timeout longer than the block time, or a waiting worker fails after almost a minute.
  7. A dead worker's task waits in the pending list; XAUTOCLAIM gives it to a live worker.
  8. Delivery is at least once: a reclaimed task ran twice, so make tasks idempotent.
  9. Count deliveries and move poison tasks to a dead-letter stream; cap streams with MAXLEN.
  10. 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.