面试知识库
中 困难

实时推送断线续传与多实例事件转发#

一句话答案#

把「推送」拆成两层:事件日志(每个会话一条 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 eid
python
  • 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 去重:

(示例假设 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")   # 让客户端重连后补发
python

7. 降级: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;完整对话和最终结果写数据库。