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 时的要点:
- 用单调时钟 + 序号或数据库自增,不要用
time.time()(同一毫秒会冲突、时钟回拨会乱序); - ID 必须对同一条流唯一,不同 run 的 ID 可以重复(回放时按 run 隔离);
- 保留一个"起始游标"(如
0-0或earliest)表示从头开始。
三、至少一次:为什么不能追求"恰好一次"
重连续播必然面临一个窗口:
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 清理订阅 |
需要显式的取消信号 |
不要把"连接断开"当成"用户取消了任务"——否则手机切后台就会杀死正在跑的长任务。
正确做法:
- 取消走独立的 API(
POST /runs/{id}/cancel),写进持久化的取消标记; - Worker 轮询该标记,协作式退出(不是 kill);
- SSE 断开时,服务端的推送停止,但运行继续;
- 客户端重连时看到的是"运行仍在进行",继续收事件。
POST /runs/{id}/cancel → 写 cancel 标记(Redis/DB)
Worker 每步检查标记 → 发现则清理并写 cancelled 事件
SSE 订阅者(含重连的)→ 收到 cancelled 事件,UI 更新
六、与发布/重启的关系
长连接在实例重启时会全部断开。所以:
- 优雅下线(drain):发布时让旧实例先停止接收新连接,等现有流结束或超时后再退出(详见构建打包与部署上线 04);
- 客户端必须能重连(这是前提,不是可选项);
- 多实例下,重连可能落到另一台实例——这只有在"事实日志外置"时才成立(回到第一节)。
⚠️ 如果事件只在进程内存里,重连落到别的实例就续不上,会表现为"有些用户能续、有些不能"——这类偶发问题极难排查,所以架构上就要杜绝进程内状态。
动手:可观察结果
| 产出 | 判断标准 |
|---|---|
| 一个事实日志 + 续播实现 | 断线后带 Last-Event-ID 重连,从断点继续,无重复无遗漏 |
| 幂等客户端 | 人为重复投递同一事件,客户端状态不变(附验证记录) |
| 过期游标处理 | 游标超出保留窗口时返回 restart 事件,客户端能回退到全量拉取 |
| 断开 vs 取消验证 | 断开连接后运行继续;调用 cancel 后运行停止(两者不混淆) |
| 多实例重连演练 | kill 一个实例,客户端重连到另一台,续播仍然正确 |
完成标志:能在压测中"每 10 秒随机掐断 5% 的连接",最终所有客户端收到的事件集合一致、顺序一致、无重复渲染。
故障注入
| 注入方式 | 观察 |
|---|---|
| 推送中途 kill 客户端进程 | 服务端是否继续运行、资源是否清理 |
重连时不带 Last-Event-ID |
是否从头重推(重复渲染) |
| 客户端处理前崩溃(已收但未应用) | 重连后是否重复投递 → 幂等是否生效 |
| 保留期设 10 秒,断线 30 秒后重连 | 是否返回 restart 而不是硬塞旧事件 |
| 事件源放在进程内存 + 多实例 | 重连到不同实例时是否续不上 |
| 把"连接断开"实现成"取消任务" | 手机切后台是否杀死正在跑的任务 |
自测题
- 为什么 SSE 重连只能做到"至少一次"?客户端该怎么应对?
- 事件 ID 应该满足哪三个性质?用 Redis Stream ID 的代价是什么?
- 游标超出保留窗口时,正确的处理方式是什么?为什么不该从头回放?
- "断开连接"与"取消运行"为什么必须分开?混在一起会造成什么后果?
- 多实例部署下,什么条件才能让"重连落到任意实例都能续"?