Appearance
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 的三个关键概念
- Stream(流):一个只追加的有序日志,类似 Kafka 的 topic。每条消息叫一个 Entry。
- Entry ID:每条消息的唯一编号,默认是
<毫秒时间戳>-<序号>(如1700000001234-0)。这个 ID 天然有序,既是 ID 也是"位置"。 - 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/Sub | Streams | List(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/目录):bashmake 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 的坑(实操容易踩)
- 发完即焚:
PUBLISH时没人订阅,消息永久丢失。排查"消息没收到"先确认订阅者是不是在发布之前就SUBSCRIBE了。 KEYS查不到 channel:channel 不是 key,不要用KEYS workspace:realtime:*找频道,用PUBSUB CHANNELS。- 订阅连接独占:
pubsub()长期占一条连接,务必用独立连接池(项目做法),否则耗尽共享池。 - 断线丢消息:Pub/Sub 没有断线补偿。需要"重连补消息"的场景,要么用 REST 拉全量(协作画布重连拉产物),要么换 Streams。
- 进订阅态不能发别的命令:
redis-cli里SUBSCRIBE后只能听;代码里给 publish 和 subscribe 配两个客户端。
9.2 Streams 的坑
- 内存会涨:Stream 持久化,不裁剪就一直堆。必须
maxlen或EXPIRE,项目两个都用了。 - 消费组忘了 ACK:
XREADGROUP读出来不XACK,消息会进 pending 列表,可能被重复消费。协作续接没用消费组(一个 message_id 一个 Stream,单消费者),所以没这问题。 - 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 拉全量兜底。
自检:能回答这几个就出师
- 用户刷新协作画布,他在断线期间别人移动光标的事件,他重连后能收到吗?为什么协作模式接受这个?
- 用户刷新聊天页,AI 正在流式回答。他重连后还能看到完整回答吗?靠什么机制?
redis_broadcaster.py为什么 publish 和 pubsub 用两个连接池?如果共用会怎样?stream_service.py的xadd同时传了maxlen和EXPIRE,为什么需要两个?- 同一个项目为什么同时用 Pub/Sub 和 Streams,而不是统一用一个?
点开看答案
- 收不到(Pub/Sub 发完即焚)。协作模式接受是因为光标是高频瞬态信息,丢几帧无所谓,且重连后会收到新的;若要看到当前画布布局,走 REST 拉全量产物兜底。
- 能。靠 Streams(
ChatResumeStreamService),AI 的每个 token 事件XADD进 Stream 持久化,用户带"上次读到的 ID"重连,XREAD从断点补齐。 pubsub()独占连接长期 listen。共用会导致订阅连接占满共享池,publish 拿不到连接报错。maxlen限制单条 Stream 的最大条数(防一条消息的事件无限堆积);EXPIRE让整个 Stream 在任务结束后自动过期清理(防 Stream key 永久残留)。两个防的是不同方向的内存泄漏。- 因为两个场景的消息可靠性需求相反:画布广播丢得起且要广播给所有人(Pub/Sub);聊天续接绝不能丢且要按断点回放(Streams)。强行统一会让简单场景变重(Pub/Sub 场景背上持久化成本),或让关键场景不可靠(Streams 场景被 Pub/Sub 的发完即焚坑)。