1. 项目概述当AI生成任务“卡在”99%相信很多在业务中集成过AI生成能力无论是文本、图像还是代码的开发者都遇到过这个让人血压飙升的场景你提交了一个耗时较长的生成任务比如生成一份长篇报告、一张高分辨率图片或者一段复杂的代码。前端进度条在稳步推进用户满怀期待地等待着。突然在进度达到90%、95%甚至99%的时候请求超时了连接断开了或者后端服务因为某个不可预知的错误崩溃了。用户看到的是一个“生成失败”的提示而你的服务器可能已经为这个任务消耗了大量的计算资源。更糟糕的是用户重试一切又得从头开始这不仅浪费资源体验也极其糟糕。这个问题的核心在于我们将AI生成这类长时间运行、非确定性的异步任务错误地用处理传统短时、同步请求的思维去架构了。传统HTTP请求-响应模型是“一锤子买卖”连接断开就意味着事务结束。但AI生成任务的生命周期远超一次HTTP连接所能维持的时间其内部状态生成了多少内容、模型参数、随机种子等是连续且宝贵的。我最近在重构一个智能文档生成系统时就深度踩了这个坑。我们的需求是用户输入一个主题AI自动生成结构完整、数据翔实的万字行业分析报告。初期方案简单粗暴前端发起一个POST请求后端调用大模型API同步等待所有内容生成完毕一次性返回。结果就是10次请求里有3次会因为网络波动或服务端超时设置而中断用户只能看到“请求超时”然后骂骂咧咧地离开。为了解决这个问题我引入了一套组合拳SSEServer-Sent Events用于实时流式输出、检查点Checkpoint用于保存中间状态、幂等Idempotency设计用于安全重试。这套方案的核心思想是将一次性的“生成”动作拆解为一个可暂停、可恢复、可观测的“状态流”。最终我们实现了即使连接在最后1%断开用户重新进入页面任务也能从断点继续并且之前已经生成的内容能立刻呈现出来体验无缝衔接。下面我就把这套方案的详细设计思路、技术选型考量、实操步骤以及填坑经验分享出来。2. 核心架构设计与思路拆解面对“AI生成到90%断了”的难题我们不能只治标比如简单调大超时时间而需要一套治本的架构。这套架构需要满足几个核心目标状态可持久化、进度可观测、操作可安全重试。我选择的三个核心技术组件正是为了分别解决这三个问题。2.1 为什么是SSE而不是WebSocket或长轮询首先解决“进度可观测”和“实时输出”的问题。我们需要一种机制让服务器能在任务执行过程中主动、持续地向客户端推送状态更新和已生成的内容片段。长轮询Long Polling技术简单但效率低下。客户端需要反复发起请求在任务完成或超时前服务器会挂起连接。对于长达数分钟的任务这会造成大量无效的连接占用和频繁的请求-响应开销不是优雅的解决方案。WebSocket功能强大支持全双工通信。但对于我们“服务器向客户端单向推送任务进度和结果片段”这个核心场景来说它有点“杀鸡用牛刀”。WebSocket需要更复杂的连接管理握手、心跳、协议升级并且对于只需要接收信息的客户端来说引入了不必要的复杂性。SSEServer-Sent Events它是基于HTTP的单向通信协议专为服务器向客户端推送事件流而设计。其优势非常契合我们的场景协议简单本质上是保持一个HTTP连接以text/event-stream格式持续发送数据流。客户端使用标准的EventSourceAPI即可连接监听无需额外库。自动重连EventSource内置了断线重连机制。当连接意外断开客户端会自动尝试重新连接并在重连时发送上一次接收到的事件ID。这为我们实现“断点续传”提供了天然支持。与HTTP生态无缝集成鉴权、缓存、代理等都可以沿用现有的HTTP基础设施部署和调试成本低。注意SSE有一个重要限制即它是文本协议传输二进制数据如图片流需要先进行Base64编码。对于文生图场景需要评估编码带来的带宽开销。在我们的文本生成场景中这是最佳选择。因此选择SSE是为了以最小的成本和复杂度实现稳定的进度推送和内容流式输出并利用其自动重连特性为恢复任务铺路。2.2 检查点Checkpoint任务状态的“存档点”检查点的概念来源于游戏和分布式计算其核心是将任务的中间状态持久化。对于AI生成任务这个“状态”远比一个简单的进度百分比复杂。一个AI文本生成任务的检查点至少需要包含任务元数据任务ID、用户ID、创建时间、使用的模型、参数如temperature, top_p。已生成的完整内容这是最重要的部分。不能只存个进度必须把已经流式输出给客户端的每一个片段都完整保存下来。模型内部状态如果可能对于某些可以控制生成过程的模型可能需要保存解码器的隐藏状态、已生成的token序列等。但大多数云端API不暴露此细节所以我们的检查点更多是“结果缓存”“续命参数”。续命参数为了能从断点继续我们需要记录“最后一句完整的话是什么”、“生成到哪个章节了”或者直接保存最后一次向模型发送的“prompt 已生成内容”作为新的上下文以便继续请求。检查点的存储介质选择内存如Redis存取速度快适合高频更新的临时状态。但服务重启数据即丢失不适合作为唯一存储。我们用它来存“实时进度”和“临时内容缓冲”。持久化数据库如PostgreSQL, MySQL作为最终状态的存储。当任务完成或客户端确认接收后将完整内容归档至此。也可用于存储检查点但频繁更新长文本字段可能对数据库有压力。对象存储如S3, MinIO存储大型内容如长文本、图片的绝佳场所。可以将每次推送的内容片段追加存储到一个文件中或者直接存储完整的生成结果。成本低扩展性好。在我们的方案中采用“Redis 数据库”的混合模式。Redis存储活跃任务的实时状态和最新内容片段数据库存储任务元数据和最终完成态。检查点的写入时机是关键。2.3 幂等Idempotency安全重试的保障幂等性意味着同一个操作执行一次或多次对系统状态产生的影响是相同的。在任务中断重试的场景下幂等设计是防止重复生成、数据混乱的“安全锁”。为什么需要幂等假设一个任务在90%时连接断开客户端自动重连SSE特性。如果没有幂等控制服务器可能错误地启动了一个全新的任务从0%开始。或者虽然尝试恢复旧任务但因为重复的请求导致旧任务被重复提交了多次造成资源浪费和结果错乱。实现任务级别的幂等通常通过**唯一的幂等键Idempotency Key**来实现。这个键通常由客户端在首次请求时生成如一个UUID并随请求头如Idempotency-Key: uuid发送。服务器端收到请求后首先以这个幂等键为主键查询是否存在未完成或已完成的任务记录。如果存在未完成的记录则直接返回该任务的当前状态和SSE连接信息实现“重连续传”。如果存在已完成的记录则直接返回最终结果。如果不存在记录则创建新任务并将幂等键与任务绑定。这样无论客户端因为网络问题发送了多少次重连请求系统都只会有一个对应的任务在执行或已执行完毕完美避免了重复劳动。3. 系统流程与核心环节实现下面我将结合代码片段以Node.js Express为例和流程图详细拆解这套系统是如何协同工作的。3.1 任务生命周期与状态流转一个具备断点续传能力的AI生成任务其状态机比简单的“进行中/完成/失败”要复杂。[客户端] --(1. 创建请求携带 Idempotency-Key)-- [服务器] [服务器] --(2. 检查幂等键)-- [Redis/DB] |- 键存在且任务未完成 - 跳转至步骤6 (重连) |- 键存在且任务已完成 - 直接返回最终结果 |- 键不存在 - 继续步骤3 [服务器] --(3. 创建任务记录状态PENDING)-- [DB] [服务器] --(4. 异步执行生成任务)-- [Worker/Async Process] [客户端] --(5. 建立SSE连接 /tasks/:id/events)-- [服务器] [服务器] --(6. 将SSE连接与任务ID绑定)-- [Redis] [Worker] --(7. 生成内容片段发布事件)-- [Redis Pub/Sub] [SSE Handler] --(8. 监听事件推送至客户端)-- [客户端] [Worker] --(9. 每生成一段更新检查点)-- [Redis] [客户端] --(10. 接收并显示片段)-- [UI] |- 连接断开 - 客户端EventSource自动重连携带Last-Event-ID |- 重连成功 - 服务器根据Last-Event-ID从检查点获取历史事件并重放然后继续推送新事件 [Worker] --(11. 生成完成更新任务状态COMPLETED持久化结果)-- [DB] [服务器] --(12. 发送[DONE]事件关闭SSE流)-- [客户端]3.2 关键代码实现解析1. 创建幂等任务端点// 使用 Express 框架 app.post(/api/generate/report, async (req, res) { const { topic, parameters } req.body; const idempotencyKey req.headers[idempotency-key]; // 客户端提供 if (!idempotencyKey) { return res.status(400).json({ error: Idempotency-Key header is required }); } // 检查幂等键 const existingTask await db.task.findUnique({ where: { idempotencyKey } }); // 情况1任务已存在且完成 if (existingTask existingTask.status COMPLETED) { return res.json({ taskId: existingTask.id, status: COMPLETED, resultUrl: existingTask.resultUrl // 直接返回结果地址 }); } // 情况2任务已存在且未完成可能是PENDING或PROCESSING if (existingTask [PENDING, PROCESSING].includes(existingTask.status)) { // 返回任务信息让客户端去连接对应的SSE流 return res.json({ taskId: existingTask.id, status: existingTask.status, streamUrl: /api/tasks/${existingTask.id}/events // SSE流地址 }); } // 情况3全新任务 const newTask await db.task.create({ data: { idempotencyKey, topic, parameters, status: PENDING, userId: req.user.id // 假设有用户认证 } }); // 异步触发任务执行放入消息队列如Bull、RabbitMQ await taskQueue.add(generate-report, { taskId: newTask.id, topic, parameters }); // 更新任务状态为处理中可选也可由Worker自己更新 await db.task.update({ where: { id: newTask.id }, data: { status: PROCESSING } }); res.json({ taskId: newTask.id, status: PROCESSING, streamUrl: /api/tasks/${newTask.id}/events }); });2. SSE事件流端点实现这是连接客户端实现实时推送和断点续传的核心。app.get(/api/tasks/:taskId/events, async (req, res) { const { taskId } req.params; const lastEventId req.headers[last-event-id]; // SSE标准重连头客户端自动发送 // 设置SSE响应头 res.writeHead(200, { Content-Type: text/event-stream, Cache-Control: no-cache, Connection: keep-alive, // 重要允许跨域如果前端分离部署 Access-Control-Allow-Origin: * }); // 发送一个初始心跳或确认连接的事件 res.write(id: ${Date.now()}\n); res.write(event: connected\n); res.write(data: ${JSON.stringify({ taskId })}\n\n); // **关键处理断点重连 - 重放历史事件** if (lastEventId) { // 根据lastEventId从Redis或DB中查询该ID之后的所有事件 const historicalEvents await redis.lrange(task:${taskId}:events, 0, -1); // 假设事件列表存在Redis let startReplaying false; for (const eventData of historicalEvents) { const event JSON.parse(eventData); // 找到最后一个已接收事件的位置然后发送之后的事件 if (event.id lastEventId) { startReplaying true; continue; } if (startReplaying) { res.write(id: ${event.id}\n); res.write(event: ${event.type}\n); res.write(data: ${JSON.stringify(event.data)}\n\n); } } } // 订阅该任务的新事件使用Redis Pub/Sub const subscriber redis.duplicate(); await subscriber.connect(); await subscriber.subscribe(task:${taskId}:stream, (message) { const event JSON.parse(message); // 按照SSE格式发送事件 res.write(id: ${event.id}\n); // 事件ID用于断点重连 res.write(event: ${event.type}\n); // 事件类型如 chunk, progress, error res.write(data: ${JSON.stringify(event.data)}\n\n); }); // 客户端关闭连接时清理订阅 req.on(close, () { subscriber.unsubscribe(task:${taskId}:stream); subscriber.quit(); console.log(Client disconnected from stream for task ${taskId}); }); });3. 后台Worker与检查点更新Worker是实际执行AI生成的部分。它需要与SSE流和检查点存储紧密配合。// 伪代码展示Worker逻辑 async function generateReportWorker(taskId, topic, parameters) { try { const taskStreamKey task:${taskId}:stream; const checkpointKey task:${taskId}:checkpoint; const eventsListKey task:${taskId}:events; // 初始化检查点从Redis加载如果不存在则新建 let checkpoint await redis.get(checkpointKey); let fullContent ; if (checkpoint) { console.log(Resuming task ${taskId} from checkpoint.); const cp JSON.parse(checkpoint); fullContent cp.fullContent || ; // 可能需要根据cp中的信息调整模型调用参数例如设置prompt为已生成内容 } else { console.log(Starting new task ${taskId}.); // 初始化检查点结构 checkpoint { fullContent: , lastSentEventId: null }; } // 模拟调用大模型API的流式接口例如OpenAI的stream: true const stream await openai.chat.completions.create({ model: gpt-4, messages: [{ role: user, content: 基于主题“${topic}”生成报告。 }], stream: true, }); let accumulatedChunk ; for await (const part of stream) { const contentDelta part.choices[0]?.delta?.content || ; if (contentDelta) { accumulatedChunk contentDelta; fullContent contentDelta; // **策略按句子或段落分割推送而不是每个token都推** // 这里简单按句号分割作为示例实际可按\n\n或固定长度分割 if (accumulatedChunk.endsWith(。) || accumulatedChunk.endsWith(.\n)) { const eventId evt_${Date.now()}_${Math.random().toString(36).substr(2, 9)}; const eventData { id: eventId, type: chunk, data: { chunk: accumulatedChunk, progress: calculateProgress(fullContent) // 估算进度 } }; // 1. 发布事件到SSE流 await redis.publish(taskStreamKey, JSON.stringify(eventData)); // 2. 将事件存入历史列表用于断线重连时重放 await redis.rpush(eventsListKey, JSON.stringify(eventData)); // 控制列表长度防止内存无限增长只保留最近N个事件或一段时间内的事件 await redis.ltrim(eventsListKey, -100, -1); // 只保留最后100个事件 // 3. 更新检查点异步或定期进行避免每次写入 // 这里采用节流方式每推送5个片段或内容长度增长超过500字符时更新一次 if (eventsListKey.length % 5 0 || fullContent.length - checkpoint.fullContent.length 500) { const newCheckpoint { fullContent: fullContent, lastSentEventId: eventId, updatedAt: new Date().toISOString() }; await redis.setex(checkpointKey, 86400, JSON.stringify(newCheckpoint)); // 设置24小时过期 } accumulatedChunk ; // 清空当前累积片段 } } } // 生成完成处理最后的累积内容如果有 if (accumulatedChunk) { // ... 同样发布事件、存储、更新检查点 ... } // 最终步骤发送完成事件持久化最终结果清理临时数据 const doneEvent { id: evt_final, type: done, data: { resultUrl: /api/results/${taskId} } }; await redis.publish(taskStreamKey, JSON.stringify(doneEvent)); await db.task.update({ where: { id: taskId }, data: { status: COMPLETED, result: fullContent, completedAt: new Date() } }); // 可选清理Redis中的临时数据检查点、事件列表或设置更短的过期时间 await redis.del(checkpointKey, eventsListKey); } catch (error) { console.error(Task ${taskId} failed:, error); // 发布错误事件 await redis.publish(taskStreamKey, JSON.stringify({ id: evt_error_${Date.now()}, type: error, data: { message: 生成任务失败, error: error.message } })); await db.task.update({ where: { id: taskId }, data: { status: FAILED, error: error.message } }); } }4. 实操避坑指南与性能优化在实际落地这套方案的过程中我遇到了不少坑也总结出一些优化经验。4.1 检查点更新的频率与粒度坑最初我每生成一个token或一个很小的片段就更新一次Redis检查点。这导致了极高的Redis IOPS在并发任务多的时候Redis成了性能瓶颈甚至影响了事件推送的实时性。解决方案采用节流Throttle和防抖Debounce思想更新检查点。时间阈值至少每2-5秒才更新一次检查点而不是实时更新。内容阈值累积生成的内容长度超过一定字符数如500字再更新。事件阈值每推送N个SSE事件后更新一次。组合使用在我们的最终方案中采用了“内容增长超过500字符或每推送5个事件”的复合条件在数据安全性和性能之间取得了很好的平衡。4.2 SSE连接管理与资源释放坑当用户离开页面时浏览器可能会关闭EventSource连接但服务器端的响应流res对象可能不会立即被Node.js垃圾回收订阅了Redis频道的连接也未释放导致内存和连接泄漏。解决方案监听req.on(close)事件这是最重要的。一旦客户端连接关闭立即取消Redis订阅并清理相关资源如上面的代码所示。设置心跳与超时在SSE流中定期发送注释行:开头的行作为心跳。可以设置一个服务器端的超时计时器如果长时间未收到客户端心跳需要客户端配合回送则主动关闭连接。使用连接池或唯一标识将SSE连接对象与任务ID关联存储在一个WeakMap或专门的连接管理器中便于在任务完成或出错时主动查找并关闭所有相关的SSE连接。4.3 历史事件重放的策略与存储坑将所有历史SSE事件无限制地存入Redis列表当生成长文档时列表可能巨大消耗大量内存且在重连时重放全部事件会导致网络流量暴增和客户端处理延迟。解决方案限制列表长度使用LRANGE和LTRIM命令只保留最近50-100个事件。因为断线重连通常发生在很短的时间窗口内不需要保存全部历史。基于检查点的增量重放这是更优雅的方案。客户端重连时发送Last-Event-ID。服务器端不需要存储所有历史事件而是从检查点中取出完整的已生成内容将其作为“第一个事件”推送给客户端。然后Worker从检查点记录的位置继续生成新内容。这样重连时只需要发送一次“全量快照”后续就是增量流。这要求客户端能处理这种“快照增量”的模式。分片存储对于超长任务可以将事件按章节或时间分片存储在不同的Key中重连时只加载最近的一个分片。4.4 幂等键的生成与生命周期管理坑客户端生成的幂等键如UUID在用户刷新页面后丢失导致无法关联到之前的任务用户被迫开始一个新任务。解决方案前端持久化幂等键在首次创建任务时将服务器返回的taskId和客户端自己生成的idempotencyKey存储在localStorage或sessionStorage中。页面刷新后优先尝试使用存储的idempotencyKey去恢复任务。服务端关联用户与会话在创建任务时不仅检查幂等键也检查当前用户是否有未完成的相同主题任务需定义“相同”的语义如主题、参数一致如果有可以提示用户是否恢复。这可以作为幂等键机制的补充。设置合理的过期时间在Redis中存储的检查点和任务状态需要设置过期时间如24小时。防止未完成的“僵尸任务”永久占用资源。数据库中的任务记录可以保留更久但状态应标记为EXPIRED或ABANDONED。4.5 前端客户端的健壮性处理前端并非只是简单监听SSE它需要处理各种边界情况。class AITaskClient { constructor(taskId, streamUrl) { this.taskId taskId; this.streamUrl streamUrl; this.eventSource null; this.retryCount 0; this.maxRetries 5; this.accumulatedData ; // 用于累积已接收的数据 this.lastEventId null; // 记录最后收到的事件ID } connect() { const headers {}; if (this.lastEventId) { headers[Last-Event-ID] this.lastEventId; // 断线重连的关键 } this.eventSource new EventSource(this.streamUrl, { headers }); this.eventSource.onmessage (event) { // 标准SSE消息event.data是字符串 const data JSON.parse(event.data); this.handleEvent(data); }; this.eventSource.addEventListener(chunk, (event) { const data JSON.parse(event.data); this.lastEventId event.lastEventId; // 更新最后事件ID this.accumulatedData data.chunk; this.updateUI(this.accumulatedData, data.progress); }); this.eventSource.addEventListener(done, (event) { const data JSON.parse(event.data); console.log(Task completed!, data.resultUrl); this.eventSource.close(); this.showFinalResult(data.resultUrl); }); this.eventSource.addEventListener(error, (event) { console.error(SSE Error:, event); // EventSource在连接失败时会自动重试但我们可以控制重试逻辑 this.eventSource.close(); this.retryCount; if (this.retryCount this.maxRetries) { setTimeout(() this.connect(), 1000 * Math.pow(2, this.retryCount)); // 指数退避重连 } else { this.showFatalError(任务连接失败请刷新页面重试。); } }); this.eventSource.onopen () { console.log(SSE连接已建立); this.retryCount 0; // 连接成功后重置重试计数 }; } handleEvent(data) { // 处理通用事件或未指定类型的事件 console.log(Received event:, data); } updateUI(content, progress) { // 更新DOM显示实时生成的内容和进度条 document.getElementById(output).innerText content; document.getElementById(progressBar).style.width ${progress}%; } } // 使用示例页面加载时尝试从localStorage恢复任务 const savedTaskId localStorage.getItem(lastTaskId); const savedIdempotencyKey localStorage.getItem(lastIdempotencyKey); if (savedTaskId savedIdempotencyKey) { // 尝试恢复任务 fetch(/api/tasks/${savedTaskId}/status, { headers: { Idempotency-Key: savedIdempotencyKey } }).then(/* ... */); }5. 扩展思考与方案变体这套“SSE检查点幂等”的模式具有很强的普适性不仅限于AI文本生成。文生图/图生图检查点可以存储生成到第几步的潜在变量Latent、去噪进度等。SSE推送的不再是文本片段而是生成过程的预览图低分辨率或模糊版本或者进度百分比。最终生成完成后再推送高清图URL。由于图片数据量大推送预览图需考虑Base64编码的带宽问题或改用WebSocket分片传输二进制数据。长视频生成/渲染检查点保存已渲染完成的帧序列或视频片段文件路径。SSE推送渲染进度和已完成的视频片段URL前端可以边下边播。复杂数据处理/ETL任务检查点保存已处理的数据批次、转换状态。SSE推送处理进度和日志信息。一个重要的权衡是状态恢复的粒度。完全无损的“热恢复”从模型内部状态继续通常需要AI服务提供方的深度支持这对于使用云端API的我们来说很难实现。我们实现的更多是一种“温恢复”或“用户无感知的冷恢复”保存所有已输出的结果当连接恢复时立刻将历史结果展示给用户同时从断点开始请求AI服务生成剩余部分。对于用户而言体验是连续的这就已经解决了核心痛点。最后这套方案引入了一定的复杂度包括需要维护Redis、消息队列、更复杂的状态管理。因此它更适合于生成耗时较长如30秒、结果价值高、用户等待耐心有限的核心场景。对于秒级完成的简单生成传统的异步轮询或Callback或许仍是更经济的选择。架构的选择永远是在用户体验、开发成本和系统复杂度之间寻找最佳平衡点。