实时推送断线续传与多实例事件转发#
一句话答案#
把「推送」拆成两层:事件日志(每个会话一条 Redis Stream,XADD 生成单调递增的 id,带长度上限和 TTL)负责不丢,实时通道(Redis Pub/Sub 广播到所有 API 实例,由持有连接的那台转给浏览器)负责快。客户端断线重连时带上最后收到的 id(SSE 由浏览器自动带
Last-Event-ID,WebSocket 要自己放在 query 或首帧里),服务端先订阅实时通道、再从日志补发、最后按 id 去重接上直播;取消指令也走广播,因为 API 不知道任务在哪台 worker 上。
核心要点
前提架构是「API 收请求并持有长连接,worker 执行 Agent」(见 Agent长任务队列与服务化)。事件在 worker 产生,连接在某台 API 上,而且连接随时会断。SSE / WebSocket 本身的对比见 流式输出与实时交互。
1. 只用 Pub/Sub 会丢什么#
Redis Pub/Sub 是 fire-and-forget:发布那一刻没有订阅者,消息就没了,也没有回放。会在三个时刻丢事件:
| 时刻 | 发生了什么 |
|---|---|
| 任务启动时 | 客户端先 POST 任务,再建连接订阅;worker 起得快,第一批事件(run_start、第一个工具调用)在订阅之前已经发完 |
| 断线重连期间 | 手机切网络、电脑休眠、网关空闲超时,重连前几秒的事件无处可去 |
| 连接迁移时 | 负载均衡把重连打到另一台 API,旧实例上的缓冲随连接一起丢了 |
所以要有一份能按位置重读的日志。实时通道和日志各管一件事:Pub/Sub 保证低延迟广播,Stream 保证可补发。
2. 事件日志:Redis Stream + 单调 id#
EVENT_LOG_MAXLEN = 500 # 只需覆盖「短时断线」期间的事件量,按单轮事件数估
EVENT_LOG_TTL = 3600 # 会话静默多久后整条清掉,按用户可能回来续看的时长定
async def emit(r, thread_id: str, event: dict) -> str:
key = f"events:{thread_id}"
eid = await r.xadd(key, {"e": json.dumps(event)}, maxlen=EVENT_LOG_MAXLEN, approximate=True) # 返回 "<毫秒>-<序号>"
await r.expire(key, EVENT_LOG_TTL) # 每次追加刷新 TTL
await r.publish("events:bus", json.dumps({"tid": thread_id, "id": eid, "e": event}))
return eidpython- id 为什么用 XADD 生成的:同一个 Stream 内 id 严格递增,即使机器时钟回拨,Redis 也会保证新 id 大于上一条。自增序号要另开 INCR,多一次往返;用时间戳自己拼 id,多个 worker 之间不保证单调。
- 比较 id 不能按字符串比:
"1700000000000-10"按字典序小于"1700000000000-9"。要拆成(毫秒, 序号)两个整数比较。 - 先写日志、再发广播:id 由 XADD 生成,本来就只能先写日志;这也保证客户端从直播里拿到的任何 id,重连时都能在日志里定位到。
- 保留策略:
MAXLEN ~ N限长(近似裁剪开销更小),加 TTL 让已结束的会话自动清理。日志只服务「短时间断线补发」,不是历史存储,最终结果和完整对话要写数据库。
3. 断线续传:带 last_event_id 重连#
SSE:服务端每帧写 id: <eid>,浏览器 EventSource 记住最后一个 id,断线后自动重连并在请求头带 Last-Event-ID;retry: 3000 可以设置重连间隔。自己用 fetch 读流时这些都没有,需要手动记录 id、手动重连并把 id 放进请求头或 query。FastAPI 0.135 起内置 fastapi.sse.EventSourceResponse,可以 yield ServerSentEvent(data=..., event=..., id=..., retry=...),重连时用 Header() 参数读 Last-Event-ID;它只负责帧格式和保活,从哪里补发仍要自己实现。
WebSocket:协议里没有这个机制,约定在重连 URL 上带 ?last_event_id=...,或在连接后的第一条消息里发送。客户端重连要用指数退避加随机抖动,避免服务端重启后所有客户端在同一秒涌回来。
服务端补发的难点是补发和直播的衔接:先补发后订阅,两步之间产生的事件会漏;先订阅后补发,同一个事件可能收两次。正确顺序是先订阅、缓冲直播,再补发,最后按 id 去重:
def id_tuple(eid: str) -> tuple[int, int]:
ms, seq = eid.split("-")
return int(ms), int(seq)
async def stream_events(r, hub, thread_id: str, last_id: str | None):
q: asyncio.Queue = asyncio.Queue(maxsize=500)
hub.subscribe(thread_id, q) # 1. 先订阅,直播事件进 q 等着
try:
cursor = last_id or "0-0"
first = await r.xrange(f"events:{thread_id}", min="-", max="+", count=1)
if last_id and (not first or id_tuple(first[0][0]) > id_tuple(last_id)):
# 保守判断:日志已过期(空),或最早一条已晚于客户端游标,中间可能被裁掉过,补不全
yield {"type": "resync"} # 前端收到后拉状态快照
last = await r.xrevrange(f"events:{thread_id}", max="+", min="-", count=1)
cursor = last[0][0] if last else last_id # 跳过补发,从最新位置接直播
# 2. 补发 cursor 之后的日志("(" 表示不含 cursor 本身,Redis 6.2+)
for eid, fields in await r.xrange(f"events:{thread_id}", min=f"({cursor}", max="+"):
yield {"id": eid, **json.loads(fields["e"])}
cursor = eid
while True: # 3. 转直播,丢掉补发里已经发过的
msg = await q.get()
if id_tuple(msg["id"]) <= id_tuple(cursor):
continue
yield {"id": msg["id"], **msg["e"]}
cursor = msg["id"]
finally:
hub.unsubscribe(thread_id, q)python(示例假设 Redis 客户端开了 decode_responses=True。)
- 补不全时怎么办:客户端离线太久,要的 id 已经被 MAXLEN 裁掉或 TTL 过期了,就发一个
resync事件,让前端调接口拉当前任务状态和已完成的消息快照,再从最新位置继续订阅。 - 任务已结束时:日志里最后一条是终态事件,补发完就可以关闭连接;日志也过期了就直接查数据库给结果。
- 另一种写法:每个连接直接对 Stream 做
XREAD BLOCK循环,从 cursor 读到最新再阻塞等待,补发和直播是同一个循环,天然没有衔接问题。代价是每个阻塞中的连接占用一个 Redis 连接,几千条长连接就是几千个 Redis 连接;Pub/Sub 方案是每个 API 实例只开一个订阅连接,在进程内按 thread_id 分发。
4. 先建连再启动任务#
即使有日志补发,也建议客户端先建立连接(拿到或生成 thread_id)、再提交任务:
- 首批事件直接走直播,不依赖补发,首字延迟更稳定。
- 客户端第一次连接时手里没有 last_event_id,如果连接晚于任务启动,它要么从头补发(需要约定「新连接默认从 0 开始读」),要么丢掉开头的事件。
两种做法都行,关键是约定清楚「没有 last_event_id 的连接从哪里读」。
5. 多实例转发:事件和取消指令#
flowchart LR
W1[worker A<br/>执行 thread t1] -->|XADD events:t1| S[(Redis Stream<br/>事件日志)]
W1 -->|PUBLISH events:bus| P((Redis Pub/Sub))
P --> A1[API 实例 1]
P --> A2[API 实例 2<br/>持有 t1 的连接]
A2 -->|WS/SSE| C[浏览器]
C -->|POST /cancel| A1
A1 -->|PUBLISH control| P
P -->|收到后按 task_id 找本地任务| W1
- 事件方向:worker 发布到公共频道,每台 API 实例都收到,只有本地持有该 thread 连接的实例转发。实例数多、事件量大时,可以按 thread 分频道(
events:{tid}),实例只订阅自己有连接的会话,代价是订阅关系要随连接增删。 - 不要自己收自己:如果 API 实例本身也会发布事件(例如入队失败、排队超时这类在 API 侧产生的事件),它同时在本地直接推送又会从频道收到一遍。给每个进程一个随机
origin,收到自己发的就跳过,否则多实例下同一事件会出现多次。 - 取消方向:用户的取消请求可能落在任意 API 实例,API 不知道任务在哪台 worker 上。做两件事:按
task_id写一个取消标记(任务还在队列里时,worker 领到后先查标记);再在控制频道广播取消指令(任务已经在跑时,持有它的 worker 执行task.cancel())。标记按 task_id 打,不按会话打,否则会误伤同一会话中刚入队的新任务。 - 不需要会话粘滞:连接可以落在任何实例,重连换实例也没关系,因为状态都在 Redis 里。这和 IM 系统「查用户在哪个网关再定向推送」的做法不同,见 IM系统设计;Agent 场景并发连接少、每条连接事件多,广播更简单。
6. 背压与慢消费者#
LLM 输出 token 的速度可能超过弱网客户端的接收速度,缓冲会在三个地方堆积:
| 位置 | 风险 | 处理 |
|---|---|---|
| Redis 给订阅连接的输出缓冲 | 某个 API 实例处理慢,Redis 端缓冲超过 client-output-buffer-limit pubsub(具体默认值见 redis.conf)后强制断开订阅连接,这台实例的所有会话同时断流 | API 的订阅循环只做「放入本地队列」,不做任何慢操作 |
| 进程内每连接队列 | 无界队列会把内存吃满 | 用有界队列;满了不阻塞订阅循环(否则一个慢客户端拖慢整台实例) |
| socket 发送缓冲 | 写阻塞 | 写超时后关闭该连接 |
有了事件日志,慢消费者可以直接处理:队列满了就断开这条连接,客户端带 last_event_id 重连后从日志补发。这比在服务端给每个慢连接无限缓冲安全得多。另一种做法是合并:多个 text_delta 合成一个再发,减少帧数;但工具调用、确认卡、终态这类事件不能丢也不能合并。
def on_bus_message(msg): # 订阅循环里调用,必须很快
for q in local_subscribers.get(msg["tid"], ()):
try:
q.put_nowait(msg)
except asyncio.QueueFull:
close_connection_of(q, reason="slow_consumer") # 让客户端重连后补发python7. 降级:Redis 出问题时#
- 日志写失败:实时推送照常,只是断线后补不回来;可以把日志写入包在熔断器里,失败时退化为只直播,并告警。
- Pub/Sub 断开:实时事件丢失,但最终结果仍在数据库里;前端在长时间收不到事件时,主动轮询任务状态接口拿进度和结果。
- 降级方向按「丢了什么」决定:少看几条进度可以接受,结果丢失或重复执行不可以。
面试回答(2分钟版)
我把推送拆成两层:事件日志负责不丢,实时通道负责快。事件日志是每个会话一条 Redis Stream,worker 每产生一个事件先 XADD,XADD 生成的 id 在同一个 Stream 里严格递增,正好当事件 id 用;Stream 设 MAXLEN 近似裁剪,每次追加刷新 TTL,已结束的会话自动清掉。写完日志再 PUBLISH 到 Pub/Sub,所有 API 实例都能收到,只有持有这个会话连接的那台转发给浏览器。之所以不能只用 Pub/Sub,是因为它发了就不管:任务启动早于订阅、断线重连期间、换实例重连,这三个时刻事件都会丢。断线续传时客户端带最后收到的 id 重连,SSE 由浏览器自动带 Last-Event-ID,WebSocket 要自己放 query 里。服务端顺序是先订阅直播缓冲起来,再用 XRANGE 从这个 id 之后补发,然后转直播并丢弃 id 小于等于已发游标的事件,id 要拆成毫秒和序号两个整数比。日志被裁掉补不全时发 resync,让前端拉快照。取消方向反过来:请求可能落在任意 API,就按 task_id 写取消标记再广播,持有任务的 worker 自己 cancel。慢消费者不给无限缓冲,有界队列满了直接断开,让它带 id 重连补发。多实例下发布方自己也订阅同一频道,要带 origin 标识跳过自己发的事件;日志写入失败时可以退化为只直播。结合项目时可以讲:Stream 的长度上限和 TTL 按什么估出来、日志写失败时降级成什么,以及怎么验证断线重连不丢不重。
追问与易错
追问方向:
- “为什么不直接用 Redis Stream 做实时推送,还要 Pub/Sub?” → 可以,每个连接
XREAD BLOCK从游标读就行,还没有衔接问题;但阻塞读占用一个 Redis 连接,连接数随在线长连接线性增长。Pub/Sub 是每个实例一个订阅连接、进程内分发,连接多时更省。 - “补发和直播之间怎么保证不丢不重?” → 先订阅并缓冲直播,再从日志补发并记录游标,最后消费缓冲时丢弃 id 不大于游标的事件;顺序反了会漏掉补发和订阅之间的事件。
- “事件 id 为什么不能用字符串比较?” → Stream id 是
毫秒-序号,序号位数不固定,字典序下-10会排在-9前面;要拆成两个整数比较。 - “客户端离线太久,要的 id 已经被裁掉了怎么办?” → 检测到日志最早一条的 id 大于客户端的 last_event_id,或者整条日志已经过期(查不到任何事件),就发 resync,前端调接口拉当前状态快照,再从最新位置订阅。
- “EventSource 断线重连会做什么?怎么让它别重连?” → 按
retry间隔重新 GET 同一个 URL,请求头带Last-Event-ID;服务端返回 204 或非text/event-stream响应时浏览器会停止重连,客户端也可以主动close()。 - “WebSocket 重连为什么要加随机抖动?” → 服务端重启时所有客户端同时断开,固定间隔重连会在同一秒同时打回来,形成重连风暴;指数退避加随机抖动把重连分散开。
- “用户在实例 1 点了取消,任务在 worker B 上,怎么送到?” → 实例 1 按 task_id 写取消标记并在控制频道广播;worker B 收到后找到本地任务调
cancel();任务还没被领走时,领取方先查标记直接跳过。 - “Redis Pub/Sub 的订阅端慢了会怎样?” → Redis 给订阅连接的输出缓冲超过
client-output-buffer-limit pubsub配置后会断开这个订阅连接,这台 API 上所有会话同时断流;所以订阅循环只做入队,不做 IO。 - “多实例下同一个事件被推了两次,可能是什么原因?” → 发布方也订阅了同一频道,本地推送一次、从频道收到又推一次;给每个进程一个 origin 标识,收到自己发的跳过;或者补发和直播衔接时没按 id 去重。
易错:
- ❌ “用了 WebSocket 就不需要续传,连接断了重新订阅就行” → 断开到重连这几秒的事件没有地方存,必须有能按位置重读的日志。
- ❌ “断开连接等于取消任务” → 任务在 worker 上跑,连接只是观察者;断线重连要能接着看,取消必须是显式指令。
- ❌ “事件日志存全量历史,用来恢复对话” → 日志只服务短时断线补发,要限长和 TTL;完整对话和最终结果写数据库。