KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA

04 · 断线重连与事件续播 — keel 龙骨

这一章回答:连接断了之后,怎么做到"用户以为没断过",而且不多收、不少收、不乱序。

这一章回答:连接断了之后,怎么做到"用户以为没断过",而且不多收、不少收、不乱序。

SSE 连接是易碎品:代理超时、手机切后台、Wi-Fi 切换、网关重连都会掐断它。浏览器会自动重连,但如果服务端不支持续播,重连的后果是从头重推——客户端会看到内容重复、状态错乱。

一、事实日志:让重连有东西可续

核心设计:

事件先落到一个持久化的「事实日志」(fact log),SSE 只是这个日志的订阅视图。
连接断了不要紧,客户端带着 Last-Event-ID 回来,从断点继续读。

这样服务端无连接状态——任何一台实例都能接住重连,因为事实在日志里而不在进程里。

事实日志的选型 适用 保留期
Redis Stream 最常用:XADD 追加,XRANGE 回放,ID 天然单调 按运行 ID 设置 TTL(如 30 分钟)
数据库事件表 需要长期审计 按业务定
Kafka 大规模、多消费者 按 topic 保留策略
进程内存 ❌ 不能用于续播(实例重启即失忆,且无法跨实例)
# 追加
await redis.xadd(f"run:{run_id}:stream", {"type": "delta", "payload": json.dumps(p)})
# 从断点读
entries = await redis.xread({key: last_id}, count=100, block=5000)

二、事件 ID 的选择(决定续播能不能简单)

理想 ID 的三个性质:单调递增、全局唯一、可直接作为"从哪里继续"的游标。

Redis Stream 的 message ID(1710000000000-0)天然满足,这就是为什么可以直接把它当 SSE 的 id:

stream_id = await redis.xadd(...)          # 形如 '1710000000123-0'
yield f"id: {stream_id}\nevent: delta\ndata: {payload}\n\n"

好处:不另造序列号,避免"日志偏移"与"SSE ID"两套坐标不一致的经典 bug。

代价:ID 格式与 Redis 强耦合(换成 Kafka 就得映射)。如果未来要换存储,建议在一开始就包一层 encode_cursor / decode_cursor。

自己实现 ID 时的要点:

三、至少一次:为什么不能追求"恰好一次"

重连续播必然面临一个窗口:

t1 服务端推出事件 101
t2 网络中断,客户端没收到
t3 客户端重连,带 Last-Event-ID: 100
t4 服务端从 101 继续 → 收到 ✅

但如果 t1 的 101 其实已经到达、只是客户端处理前崩溃了:
t3 带 Last-Event-ID: 100 → 101 会被再发一次 → 重复

结论:续播只能保证"至少一次"(at-least-once),重复交给客户端处理。追求"恰好一次"需要分布式事务,不值得。

所以客户端必须幂等,三种做法:

做法 适用 说明
按 ID 去重 通用 记住已处理的最大 ID,小于等于它的直接丢弃
按序号拼接 token 流 文本按 index 追加,重复的直接覆盖
按状态幂等 状态类事件 事件带完整状态而非增量,重复应用结果一致
let lastSeq = -1;
function handle(evt) {
  if (evt.seq <= lastSeq) return;      // 重复,丢弃
  lastSeq = evt.seq;
  apply(evt);
}

设计建议:给每个事件带上序号 + 完整状态(而不是纯增量),幂等会变得非常容易。这也是很多 Agent 平台选择"全量 projection + 事件通知"混合模式的原因——SSE 传的是"变了"的信号与结构化事件,而不是让用户自己去拼增量。

四、回放边界:从多早开始续?

客户端带 Last-Event-ID: X
  ├─ X 仍在保留窗口内 → 从 X 之后继续 ✅
  └─ X 太老(已被裁剪)  → 两种选择:
       ① 返回 204/终止事件,让客户端走「全量拉取」再重建状态
       ② 从头回放(可能造成大量重复,需客户端强幂等)

推荐 ①:明确告诉客户端"续不上了,请重新拉取当前状态",而不是硬塞一堆旧事件。

if cursor_too_old(last_event_id):
    yield "event: restart\ndata: {\"reason\":\"expired\"}\n\n"
    return   # 客户端收到后去调 GET /runs/{id} 拿全量状态

保留期的取值:略大于"最长可能断线时间 + 重试时间"。如果业务上允许断线 10 分钟还能续,保留 30 分钟比较稳妥。用 Redis 的话给 key 设 TTL 即可。

五、取消:断开 ≠ 取消

这是最容易混淆的一对概念:

断开连接 取消运行
触发 网络/客户端走了 用户主动点"停止"
后果 运行继续(服务端还在跑) 运行应停止
实现 finally 清理订阅 需要显式的取消信号

不要把"连接断开"当成"用户取消了任务"——否则手机切后台就会杀死正在跑的长任务。

正确做法:

  1. 取消走独立的 API(POST /runs/{id}/cancel),写进持久化的取消标记;
  2. Worker 轮询该标记,协作式退出(不是 kill);
  3. SSE 断开时,服务端的推送停止,但运行继续;
  4. 客户端重连时看到的是"运行仍在进行",继续收事件。
POST /runs/{id}/cancel  →  写 cancel 标记(Redis/DB)
Worker 每步检查标记    →  发现则清理并写 cancelled 事件
SSE 订阅者(含重连的)→  收到 cancelled 事件,UI 更新

六、与发布/重启的关系

长连接在实例重启时会全部断开。所以:

⚠️ 如果事件只在进程内存里,重连落到别的实例就续不上,会表现为"有些用户能续、有些不能"——这类偶发问题极难排查,所以架构上就要杜绝进程内状态。


动手:可观察结果

产出 判断标准
一个事实日志 + 续播实现 断线后带 Last-Event-ID 重连,从断点继续,无重复无遗漏
幂等客户端 人为重复投递同一事件,客户端状态不变(附验证记录)
过期游标处理 游标超出保留窗口时返回 restart 事件,客户端能回退到全量拉取
断开 vs 取消验证 断开连接后运行继续;调用 cancel 后运行停止(两者不混淆)
多实例重连演练 kill 一个实例,客户端重连到另一台,续播仍然正确

完成标志:能在压测中"每 10 秒随机掐断 5% 的连接",最终所有客户端收到的事件集合一致、顺序一致、无重复渲染。

故障注入

注入方式 观察
推送中途 kill 客户端进程 服务端是否继续运行、资源是否清理
重连时不带 Last-Event-ID 是否从头重推(重复渲染)
客户端处理前崩溃(已收但未应用) 重连后是否重复投递 → 幂等是否生效
保留期设 10 秒,断线 30 秒后重连 是否返回 restart 而不是硬塞旧事件
事件源放在进程内存 + 多实例 重连到不同实例时是否续不上
把"连接断开"实现成"取消任务" 手机切后台是否杀死正在跑的任务

自测题

  1. 为什么 SSE 重连只能做到"至少一次"?客户端该怎么应对?
  2. 事件 ID 应该满足哪三个性质?用 Redis Stream ID 的代价是什么?
  3. 游标超出保留窗口时,正确的处理方式是什么?为什么不该从头回放?
  4. "断开连接"与"取消运行"为什么必须分开?混在一起会造成什么后果?
  5. 多实例部署下,什么条件才能让"重连落到任意实例都能续"?

进入 keel 阅读