SSE断联重试
SSE断联重试
一、背景
在大模型 ChatBot、AI Agent 等应用中,为了提供类 ChatGPT 的实时体验,服务端通常采用 SSE(Server-Sent Events)将模型生成的 Token 逐字推送到前端。典型流程为:用户提问 → 后端调用 LLM → LLM 持续生成 Token → SSE 实时推送 → 前端打字机展示。
然而,SSE 基于 HTTP 长连接,无法保证永久稳定。用户刷新页面、网络抖动、切换 WiFi、关闭标签页等行为都会导致连接断开。若不处理,用户将丢失未完成的生成结果,体验大幅下降。
核心问题在于:`SSE 断开后,如何让用户重连并继续接收未完成的输出?
二、误区:将 SSE 连接等同于 LLM 任务
常见的初级实现将 SSE 连接与 LLM 调用直接绑定:@GetMapping("/chat") public Flux<String> chat() { return llm.stream(); }。这种模式下,一旦 SSE 断开,HTTP 请求被取消,LLM 生成任务也随之终止。结果是已生成内容丢失且无法恢复,已消耗的计算资源被浪费,用户必须重新提问从头生成。
因此,生产环境必须将 SSE 连接 与 LLM 生成任务 解耦:SSE 仅作为事件订阅通道,而非任务执行载体。
三、方案设计
核心思想是将 LLM 生成的每个输出块(chunk)持久化到可靠的事件流中,同时实时推送给当前 SSE 连接。断开重连后,新连接从事件流中读取未消费的部分,实现“断点续传”。
整体架构为:用户 → SSE 连接 → Chat Gateway,Gateway 同时对接事件流存储(如 Redis Stream / Kafka)和 LLM Worker。LLM Worker 独立运行生成任务,每生成一个 chunk,同时写入事件流并推送至当前 SSE。事件流存储按 taskId 分区存储有序事件,支持按偏移量读取。Chat Gateway 管理 SSE 连接,处理重连请求,从事件流中拉取历史并实时订阅新事件。
四、正常流程与重连恢复
正常生成流程:LLM 生成 "Redis"、"是"、"数据库" 三个 chunk。首先写入事件流:seq=1:Redis、seq=2:是、seq=3:数据库(按 taskId 存储);然后推送 SSE 给前端,用户看到打字机效果。写入与推送同步进行,用户感知不到额外延迟。
刷新后恢复流程:假设用户刷新前已收到 seq=1 和 seq=2,前端保存了 taskId 和 lastEventId = 2。刷新后发起新 SSE 请求,携带 Last-Event-ID: 2。后端从事件流中读取 seq > 2 的所有事件(seq=3:数据库 及后续),推送给新连接。若 LLM 仍在运行,后续新生成的事件也会继续推送。用户看到的只是页面闪了一下,内容无缝衔接。
五、关键技术细节
事件粒度与存储策略:不保存每个 Token(粒度过细,存储和回放开销大),实践中将 Token 聚合成语义块(Chunk),如按时间窗或完整词语合并。生成过程中的事件保存在 Redis Stream(设置 TTL),生成完成后将最终完整答案存入 MySQL,并清理临时事件。
任务取消策略:SSE 断开后是否立即取消 LLM 任务取决于业务场景。普通聊天中用户刷新可能只是为了看完整结果,可让 LLM 继续生成,保留结果供重连取用;Agent 任务(涉及工具调用、多轮协作)成本较高,可设置超时时间(如 30 秒),若未重连则主动取消,释放资源。
消费模式:每个 SSE 连接通过阻塞读取(XREAD BLOCK)订阅事件流,有新事件时立即返回,无新事件时挂起,不产生无效轮询。
六、代码示意(简化)
以下使用 Spring WebFlux + Redis Stream 演示核心逻辑。
生成端:LLM Worker 每生成一个 chunk 调用 onChunk 方法,写入 Redis Stream 并推送。
public void onChunk(String taskId, String chunk, long seq) {
redisTemplate.opsForStream().add("chat:task:" + taskId, Map.of("seq", seq, "data", chunk));
sink.next(chunk);
}重连端:Controller 接收 taskId 和可选的 Last-Event-ID,从 Stream 中读取大于该序号的事件,并继续监听新事件。
@GetMapping(value = "/stream/{taskId}", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
public Flux<String> stream(@PathVariable String taskId,
@RequestHeader(value = "Last-Event-ID", required = false) Long lastId) {
long startSeq = lastId != null ? lastId + 1 : 0;
return redisTemplate.opsForStream()
.read(StreamOffset.fromStart("chat:task:" + taskId))
.filter(record -> (Long) record.getValue().get("seq") >= startSeq)
.map(record -> (String) record.getValue().get("data"))
.concatWith(/* 继续监听新事件 */);
}七、总结
大模型流式输出中的 SSE 断连问题,本质是长时间运行任务与瞬时客户端连接的矛盾。解决思路不是让连接更稳定,而是让生成任务不再依赖单一连接。
生产级方案的核心原则:LLM 任务与 SSE 连接解耦;生成过程以事件流形式持久化;SSE 仅作为事件订阅通道;断线后通过 taskId + offset 恢复;生成完成后持久化最终结果。
正确的关系链是 LLM → 事件流(持久化) → SSE(订阅) → 用户,而非 LLM → SSE → 用户。这也是当前大模型应用、Agent 系统实现可靠流式输出的标准范式。
一句话记住:用可回放的事件流取代易断的连接,让用户刷新永不丢失。