Skip to content

Redis Pub/Sub 进阶:Streams 对比 + 项目读码 ​

配套入门篇:docs/协作模式/2026-07-04-redis-pubsub-入门教程.md(读完再看这篇)。

本篇解决一个问题:Pub/Sub 发完即焚,消息丢了怎么办? 答案是 Redis Streams。讲完原理,带你逐行读协作模式项目里的两段真实代码——一段用 Pub/Sub 做实时广播,一段用 Streams 做聊天断线续接,看看为什么同一个项目里两者都用。


1. 先回到痛点:Pub/Sub 哪里不够用? ​

入门篇我们建立了一个直觉:Pub/Sub 是电台直播,发完即焚。这在协作画布光标、在线状态这种"丢了无所谓"的场景很合适。

但假设有这样一个场景:

用户在跟 AI 聊天,AI 正在流式输出回答(一个字一个字蹦)。这时用户刷新了网页。按 Pub/Sub 的逻辑,刷新期间产生的所有流式 token 全部丢失,用户重连后只能看到一片空白,因为消息发完即焚,根本没存。

这个场景下我们想要的是:

  • 消息要存起来(持久化),不是发完即焚。
  • 消费者(浏览器)断线重连后,能从断点继续读之前错过的消息。
  • 最好每条消息有个顺序号,方便定位"我读到哪了"。

这就是 Redis Streams 干的事。


2. 费曼类比:Streams 是"带时间轴的录音带" ​

如果说:

  • Pub/Sub = 现场直播:你只能听"正在播的",错过了回不去。
  • List 做队列 = 一沓一次性便签:消费者拿走一张就撕掉(LPUSH/RPOP),别人看不到。
  • Streams = 录音带(带时间戳的有序日志):每句话都录下来标上时间戳,谁都能随时倒带从任意位置重听,听过的还能标记"我听到这了"。

Streams 的三个关键概念 ​

  1. Stream(流):一个只追加的有序日志,类似 Kafka 的 topic。每条消息叫一个 Entry。
  2. Entry ID:每条消息的唯一编号,默认是 <毫秒时间戳>-<序号>(如 1700000001234-0)。这个 ID 天然有序,既是 ID 也是"位置"。
  3. Consumer Group(消费组):多个消费者组成一组,共同消费一个 Stream,每条消息只被组内一个消费者处理,处理完确认(ACK)。没确认的消息可以重新投递——这就是 Pub/Sub 没有的"可靠性"。

关键区别:Streams 里每条消息都"待在那里",直到你主动裁剪(XTRIM/maxlen)或设过期。消费者掉线了,上线后可以用记住的 ID 回去把错过的消息全读回来。


3. Streams 核心命令 ​

命令谁用干啥对应入门篇的
XADD stream * 字段 值生产者往流里追加一条(* 让 Redis 自动生成 ID)PUBLISH
XLEN stream任意查流里有几条(Pub/Sub 没有对应物)
XRANGE stream - +任意按 ID 范围读(- 最老,+ 最新)(没有,Pub/Sub 不能倒带)
XREAD COUNT n STREAMS stream id消费者从某 ID 之后读;$ 表示只读未来的新消息SUBSCRIBE
XREAD BLOCK 0 STREAMS stream $消费者阻塞等待新消息(类似 SUBSCRIBE 的 listen)SUBSCRIBE 阻塞
XGROUP CREATE ...初始化创建消费组(无)
XREADGROUP GROUP 组 消费者 ...消费者消费组读取(自动分配,保证只给一个消费者)(无)
XACK stream 组 id消费者确认处理完(否则可被 XPENDING/XCLAIM 重投)(无)
XTRIM stream MAXLEN 1000维护裁剪,只保留最近 N 条(无,Pub/Sub 不存)

记忆窍门:带 X 开头的命令都是 Stream 专属。最常用的就 XADD(写)、XRANGE/XREAD(读)、XLEN(看长度)。


4. 实操:用 redis-cli 手敲 Streams ​

接着入门篇的 Redis(同样 docker run --rm -p 6379:6379 redis:7)。

4.1 追加消息(XADD) ​

127.0.0.1:6379> XADD chat.msg * user alice text "你好"
"1700000001234-0"     ← Redis 自动生成的 ID(时间戳-序号)

127.0.0.1:6379> XADD chat.msg * user bob text "在吗"
"1700000001567-0"

127.0.0.1:6379> XADD chat.msg * user alice text "今天天气不错"
"1700000001890-0"

注意 * —— 让 Redis 自动按时间戳生成 ID。你也可以手写 ID,但通常没必要。

4.2 看长度和全部内容(XLEN / XRANGE) ​

127.0.0.1:6379> XLEN chat.msg
(integer) 3

127.0.0.1:6379> XRANGE chat.msg - +
1) 1) "1700000001234-0"           ← 消息 ID
   2) 1) "user"                   ← 字段
      2) "alice"
      3) "text"
      4) "你好"
2) 1) "170000000567-0"
   2) ...

这一步就是 Pub/Sub 做不到的事:消息存下来了,你随时 XRANGE 倒带重读。

4.3 从某个位置之后读(XREAD) ​

XREAD 是 Streams 的"订阅",但它能指定从哪个 ID 之后开始:

# 只读"从现在开始"的新消息($ = 当前最新 ID)
127.0.0.1:6379> XREAD COUNT 10 STREAMS chat.msg $
(nil)     ← 此刻没有比 $ 更新的消息

# 从指定 ID 之后读(模拟"我断线重连,从这之后补消息")
127.0.0.1:6379> XREAD STREAMS chat.msg 1700000001234-0
1) 1) "chat.msg"
   2) 1) 1) "1700000001567-0"     ← 跳过了 1234 那条,只给之后的
         2) ...

💡 这正是协作模式"聊天续接"的核心:浏览器刷新后,带着"我上次读到的最后一条 event ID"请求后端,后端用 XREAD stream 上次ID 把断线期间错过的流式事件全补回来。Pub/Sub 根本做不到。

4.4 阻塞等待新消息(XREAD BLOCK) ​

# 阻塞最多 10 秒等新消息(类似 SUBSCRIBE 的 listen)
127.0.0.1:6379> XREAD BLOCK 10000 STREAMS chat.msg $

此时在另一个终端 XADD chat.msg * user alice text "来了",这个终端会立刻收到。

4.5 裁剪(XTRIM / maxlen) ​

消息不能无限堆积,需要裁剪:

127.0.0.1:6379> XTRIM chat.msg MAXLEN 1000    # 只保留最近 1000 条

项目代码里用 xadd(..., maxlen=STREAM_MAX_LEN, approximate=True) 在写入时顺手裁剪。approximate=True 是"近似裁剪"(性能更好,可能多留几十条),生产推荐。


5. Pub/Sub vs Streams vs List:一张表选型 ​

维度Pub/SubStreamsList(LPUSH/RPOP)
持久化❌ 发完即焚✅ 持久化(直到裁剪/过期)✅ 持久化(直到取出)
断线补消息❌ 不可能✅ 按 ID 回读⚠️ 取走即消失,多人难协作
消费模式广播(所有订阅者都收)广播(XREAD)或消费组(竞争)竞争(RPOP 抢)
顺序保证单频道有序全局有序(按 Entry ID)队列有序
消息确认(ACK)❌✅ 消费组 XACK❌
失败重投❌✅ XPENDING/XCLAIM需自己实现
多消费者都能收全量✅✅(各自 XREAD 不用组)❌(一个抢走就没了)
内存占用极低(不存)取决于 maxlen/裁剪取决于消费速度
典型场景实时广播、在场态任务队列、事件溯源、断线续接简单任务队列

决策树 ​


6. 项目读码:为什么协作模式两个都用? ​

这是本篇最有价值的部分。同一个 backend/ 里,Pub/Sub 和 Streams 各管一摊,它们的选择不是随意的,而是由场景的消息可靠性需求决定。

6.1 两个场景对照 ​

场景用什么为什么
协作画布实时广播(光标移动、产物新增、有人上线)Pub/Sub丢一两个光标事件无所谓,下一秒还有新的;且要广播给所有在线成员
聊天流式输出续接(刷新页面后继续看 AI 的回答)Streams流式 token 绝不能丢,用户刷新后必须能补回完整回答

7. 逐行读码:Pub/Sub 实战(画布广播) ​

文件:backend/services/workspace_realtime/redis_broadcaster.py。这就是入门篇你写的 Python demo 的生产加强版。

7.1 发布(就一行) ​

python
# backend/services/workspace_realtime/redis_broadcaster.py
_CHANNEL_PREFIX = "workspace:realtime"

def _channel(workspace_id: str) -> str:
    return f"{_CHANNEL_PREFIX}:{workspace_id}"   # 频道名:workspace:realtime:{workspace_id}

async def publish(self, workspace_id: str, event: dict[str, Any]) -> None:
    # 就是入门篇 publisher.py 里 r.publish(...) 的异步版
    await self._redis.publish(_channel(workspace_id), json.dumps(event, ensure_ascii=False))

每个协作项目(workspace)对应一个独立频道。画布操作 → 序列化成 JSON → 往这个频道喊。和入门篇你的 demo 一模一样,只是频道名按 workspace 划分。

7.2 订阅循环(就是入门篇的 listen,加了重连) ​

python
async def _pubsub_loop(self, workspace_id: str) -> None:
    while True:                        # 外层:重连循环,Redis 抖动后退避重订阅
        try:
            async with self._pubsub_redis.pubsub(ignore_subscribe_messages=True) as pubsub:
                await pubsub.subscribe(_channel(workspace_id))   # 入门篇的 subscribe
                while True:
                    message = await pubsub.get_message(timeout=1.0)  # 入门篇的 listen
                    if message is None or message.get("type") != "message":
                        continue
                    event = self._decode(message.get("data"))
                    if event is None:
                        continue
                    # 把消息分发给本实例注册的本地 handler(转发给各 WS 连接)
                    for handler in list(self._subscriptions.get(workspace_id, set())):
                        result = handler(workspace_id, event)
                        if asyncio.iscoroutine(result):
                            await result
        except asyncio.CancelledError:
            raise                      # 正常取消(退订),退出
        except Exception:
            await asyncio.sleep(_RECONNECT_BACKOFF_SECONDS)   # 异常退避重连

比你的 demo 多的两点工程化:①while True 重连退避(入门篇提过 Pub/Sub 断线丢消息,这里靠"重连"恢复后续消息,断线期间的依然会丢,但协作画布可接受);②收到后分发给本地连接表(ConnectionRegistry),因为一个 Redis channel 收到的消息要扇出给本实例上同 workspace 的所有 WebSocket。

7.3 为什么要独立连接池(回顾入门篇 4.1) ​

python
def _make_pubsub_redis() -> aioredis.Redis:
    # pubsub 独立池(max=64),与 publish 共享池(max=50)隔离
    return aioredis.from_url(config.REDIS_URL, decode_responses=True, max_connections=64)

入门篇 4.1 讲过原因:pubsub() 会独占连接长期 listen,共池会耗尽 publish 的连接。看到这里你应该完全懂了。

7.4 配套调试工具:redis_debug.py ​

文件:backend/scripts/redis_debug.py。这是排查 Pub/Sub 的"眼睛"——因为 Pub/Sub 发完即焚、KEYS 查不到,不实时监听根本看不到消息。

python
# backend/scripts/redis_debug.py(精简)
def cmd_sub(workspace_id):
    channel = _channel(workspace_id)
    client = redis.Redis.from_url(config.REDIS_URL, decode_responses=True)
    pubsub = client.pubsub(ignore_subscribe_messages=True)
    pubsub.subscribe(channel)              # 调台
    for message in pubsub.listen():        # 持续听
        obj = json.loads(message["data"])
        print(f"[{now}] {obj.get('type')} ...")

用法(在 backend/ 目录):

bash
make redis-sub ARGS=<workspace_id>        # 实时监听某 workspace 的协作消息流
make redis-presence ARGS=<workspace_id>   # 查在线成员 + 光标快照(presence hash)
make redis-channels                        # 列出当前所有正在协作的 channel

这个工具就是入门篇你写的 subscriber.py 的项目落地版,加上了 JSON 解析和频道名约定。下次协作画布消息不灵,直接 make redis-sub 监听一眼就知道后端到底有没有发出消息。


8. 逐行读码:Streams 实战(聊天续接) ​

文件:backend/services/core/resume/stream_service.py,类 ChatResumeStreamService。

业务背景:AI 回答用户问题时,是流式输出的(token 一个个蹦)。如果用户中途刷新页面,传统做法是回答丢了重头再来。这里用 Streams 把每个 token 事件录进 Stream,用户刷新后用 message_id 找到 Stream,从断点把剩余 token 补回来,实现无缝续接。

8.1 写入:XADD(录进录音带) ​

python
# backend/services/core/resume/stream_service.py
async def append_event(self, message_id: int, event_type: str, payload: dict[str, Any]) -> str | None:
    """向 Redis Stream 追加一条可恢复事件。"""
    return await get_redis().xadd(
        redis_keys.stream(message_id),       # stream key,如 chat:stream:{message_id}
        {
            "type": event_type,              # 事件类型(token/chunk/done...)
            "payload": json.dumps(payload, ensure_ascii=False),
            "created_at": format_shanghai_iso(now_shanghai_naive()) or "",
        },
        maxlen=STREAM_MAX_LEN,               # 裁剪:只保留最近 N 条,防无限膨胀
        approximate=True,                    # 近似裁剪(性能好,可能多留几十条)
    )
    # 返回值 = 这条事件的 Entry ID(时间戳-序号),前端拿着它当"读到哪里了"的书签

对照入门篇:

  • xadd ≈ PUBLISH,但消息存下来了,且返回一个 Entry ID 当书签。
  • maxlen + approximate = 自动裁剪,对应你 XTRIM 那一步,只是写入时顺手做了。
  • 每个 AI 消息一个独立 Stream(chat:stream:{message_id}),消息结束清理。

8.2 读全部:XRANGE(倒带) ​

python
async def read_events(self, message_id: int, start="-", end="+") -> list[dict[str, Any]]:
    """读取 Stream 事件并反序列化。"""
    rows = await get_redis().xrange(redis_keys.stream(message_id), min=start, max=end)
    return [self._decode_event(event_id, fields) for event_id, fields in rows]

min="-", max="+" 就是你 cli 里敲的 XRANGE stream - +(全量倒带)。用户刷新后,前端可能先调这个把已有的事件全读出来,渲染出已生成的部分回答。

8.3 读断点之后:XREAD(续接核心) ​

python
async def read_new_events(self, message_id: int, after: str, block_ms: int = 15000):
    """阻塞读取指定事件 ID 之后的新事件。"""
    rows = await get_redis().xread(
        {redis_keys.stream(message_id): after or "$"},   # after = 上次读到的 ID;"$" = 只读未来的
        block=block_ms,          # 最多阻塞 15 秒等新消息(类似 XREAD BLOCK)
        count=50,
    )

这一段是整个续接的灵魂,对照你在 cli 里敲的 XREAD BLOCK 10000 STREAMS chat.msg $:

  • after 参数:前端传"我上次读到的最后一条 Entry ID"。Streams 会只返回这个 ID 之后的事件——也就是断线期间错过的 token,精准补齐。
  • "$":如果前端没传(第一次连),表示"从最新位置开始,只读未来的"。
  • block=15000:没有新消息时阻塞等 15 秒,而不是立刻返回空。这就是流式输出的"长轮询"效果——AI 还在生成时,后端阻塞等下一个 token,有了立刻返回。

⚠️ 对比 Pub/Sub:这里如果用 Pub/Sub,用户刷新的那几秒,流式 token 全部丢失,根本补不回来。Streams 的"持久化 + 按 ID 回读"是续接能成立的根本。

8.4 过期清理:EXPIRE(录音带别永久占地方) ​

python
# Stream key 会在写入时设过期(见文件里 ACTIVE_TTL_SECONDS 相关逻辑)
await get_redis().expire(redis_keys.stream(message_id), ACTIVE_TTL_SECONDS)

Stream 不像 Pub/Sub 用完即焚,它一直占内存。所以:

  • 写入时 EXPIRE 设一个 TTL(比如几分钟),消息结束/用户离开一段时间后自动清理。
  • 配合 maxlen 双保险:maxlen 防单条消息过长,EXPIRE 防整个 Stream 永久残留。

8.5 串起来:续接的完整时序 ​


9. 坑点与小结 ​

9.1 Pub/Sub 的坑(实操容易踩) ​

  1. 发完即焚:PUBLISH 时没人订阅,消息永久丢失。排查"消息没收到"先确认订阅者是不是在发布之前就 SUBSCRIBE 了。
  2. KEYS 查不到 channel:channel 不是 key,不要用 KEYS workspace:realtime:* 找频道,用 PUBSUB CHANNELS。
  3. 订阅连接独占:pubsub() 长期占一条连接,务必用独立连接池(项目做法),否则耗尽共享池。
  4. 断线丢消息:Pub/Sub 没有断线补偿。需要"重连补消息"的场景,要么用 REST 拉全量(协作画布重连拉产物),要么换 Streams。
  5. 进订阅态不能发别的命令:redis-cli 里 SUBSCRIBE 后只能听;代码里给 publish 和 subscribe 配两个客户端。

9.2 Streams 的坑 ​

  1. 内存会涨:Stream 持久化,不裁剪就一直堆。必须 maxlen 或 EXPIRE,项目两个都用了。
  2. 消费组忘了 ACK:XREADGROUP 读出来不 XACK,消息会进 pending 列表,可能被重复消费。协作续接没用消费组(一个 message_id 一个 Stream,单消费者),所以没这问题。
  3. Entry ID 是字符串:比较大小按字典序,但因为格式是 时间戳-序号,字典序 == 时间序。传 after 时用前端拿到的原始 ID 字符串,别自己拼。

9.3 选型一句话 ​

丢得起的实时广播 → Pub/Sub;丢不得的消息流 → Streams。

协作模式两个都用,正是因为画布广播"丢得起"、聊天续接"丢不得",各取所需。


10. 回看:你现在应该能完全读懂这两张图 ​

回头翻 docs/协作模式/2026-07-04-后端实现总结-agent产物同步与实时协同.md:

  • 第 3.5 节"广播分层"图:单实例 ConnectionRegistry + 跨实例 RedisBroadcaster 扇出 —— 现在你知道为什么用 Pub/Sub(光标/产物通知丢得起、要广播给所有在线成员)。
  • redis_broadcaster.py 的 _pubsub_loop 重连退避 —— 你知道它只能保证"重连后的消息",断线瞬间的依然会丢,但画布靠 REST 拉全量兜底。

自检:能回答这几个就出师 ​

  1. 用户刷新协作画布,他在断线期间别人移动光标的事件,他重连后能收到吗?为什么协作模式接受这个?
  2. 用户刷新聊天页,AI 正在流式回答。他重连后还能看到完整回答吗?靠什么机制?
  3. redis_broadcaster.py 为什么 publish 和 pubsub 用两个连接池?如果共用会怎样?
  4. stream_service.py 的 xadd 同时传了 maxlen 和 EXPIRE,为什么需要两个?
  5. 同一个项目为什么同时用 Pub/Sub 和 Streams,而不是统一用一个?
点开看答案
  1. 收不到(Pub/Sub 发完即焚)。协作模式接受是因为光标是高频瞬态信息,丢几帧无所谓,且重连后会收到新的;若要看到当前画布布局,走 REST 拉全量产物兜底。
  2. 能。靠 Streams(ChatResumeStreamService),AI 的每个 token 事件 XADD 进 Stream 持久化,用户带"上次读到的 ID"重连,XREAD 从断点补齐。
  3. pubsub() 独占连接长期 listen。共用会导致订阅连接占满共享池,publish 拿不到连接报错。
  4. maxlen 限制单条 Stream 的最大条数(防一条消息的事件无限堆积);EXPIRE 让整个 Stream 在任务结束后自动过期清理(防 Stream key 永久残留)。两个防的是不同方向的内存泄漏。
  5. 因为两个场景的消息可靠性需求相反:画布广播丢得起且要广播给所有人(Pub/Sub);聊天续接绝不能丢且要按断点回放(Streams)。强行统一会让简单场景变重(Pub/Sub 场景背上持久化成本),或让关键场景不可靠(Streams 场景被 Pub/Sub 的发完即焚坑)。