ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Agent流式输出管道实战:SSE、StreamChunk与断流恢复

2026/10/4 14:00:49 拓冰建站 浏览量
Agent流式输出管道实战:SSE、StreamChunk与断流恢复 1. 流式输出为什么是 Agent 体验的分水岭做过 Agent 项目的人大概都有这个体会模型能力再强如果前端等个十几秒才一次性把整段回复吐出来用户的心理感受就是卡死了。而一旦把流式输出打通同样的模型、同样的响应时间用户会觉得这玩意儿在思考、在打字体验直接上一个台阶。这就是我为什么在 DeepSeek-Harness 这个项目里把流式输出管道单独拎出来做一章的原因——它不是锦上添花的功能而是 Agent 产品能不能用的底线。这一章要聊的核心是把模型侧返回的StreamChunk一路送到 UI 上渲染出来的完整链路。中间会经过 SSE 传输、事件封装、前端解析、增量渲染这几个关键环节每一环都有坑。关键词里出现的StreamChunk、流式输出、SSE、通知sse、封装sse流式接口调用逻辑基本就是这条管道的全部关键词。适合谁看如果你正在用 DeepSeek 或者别的模型 API 搭 Agent前端用 Vue 或者 React后端用 Python并且被流式断流idle timeout消息拼接错乱这些问题折磨过那这篇就是写给你的。我先把结论摆前面流式输出管道最难的不是怎么把 token 传过去而是边界处理——什么时候算一个 chunk 结束、断流了怎么恢复、多个工具调用怎么和文本流交织、前端怎么保证渲染不闪烁。这些细节决定了你的 Agent 是能用还是好用。2. 整体架构设计与方案选型思路2.1 为什么是 SSE 而不是 WebSocket很多人第一反应是流式输出就该用 WebSocket双向、实时、听起来就专业。但我在 Harness 里最终选了 SSE理由很实在。Agent 的流式输出本质上是单向的服务端持续推 token客户端只管收。这种场景下 WebSocket 的双向能力是浪费的反而带来额外的连接管理成本——心跳、重连、状态机每一样都要写代码维护。SSE 基于普通 HTTP天然支持自动重连浏览器内置的EventSource就有retry机制服务端就是一个普通的流式响应部署时 Nginx、网关这些中间件对它的兼容性也更好。更关键的一点SSE 走的是标准 HTTP 语义意味着你可以复用现有的鉴权、限流、日志中间件。WebSocket 升级握手那一步经常和网关打架我踩过不止一次。当然 SSE 也有短板比如浏览器对同域名并发连接数有限制HTTP/1.1 下大约 6 个但 Agent 场景一般一个会话一条流问题不大。提示如果你的 Agent 需要客户端在流式过程中频繁回传中断信号、调整参数那 WebSocket 更合适。但如果只是服务端推、客户端收SSE 是更省心的选择。2.2 StreamChunk 的数据结构设计整条管道的基石是StreamChunk这个数据结构。它不能只是一个字符串因为 Agent 的输出远比纯文本复杂——有正文、有思考过程、有工具调用、有引用来源、有结束标记。我把它设计成一个带type字段的联合结构核心字段大致是这样class StreamChunk: type: str # text / reasoning / tool_call / tool_result / done / error content: str # 增量文本内容 index: int # 序号用于排序和去重 message_id: str # 所属消息 ID metadata: dict # 工具名、参数、引用等附加信息为什么要带index因为网络传输不保证顺序尤其是经过多层代理之后。前端拿到乱序的 chunk 如果直接拼接文本就会错位。有了单调递增的index前端可以做个简单的缓冲排序保证渲染顺序正确。message_id则是为了支持一条会话里多条消息并发流式——比如用户连续发了两条两条的流可能交织返回靠message_id区分。type字段是整个设计的灵魂。它让前端可以针对不同类型做不同渲染text走正文气泡reasoning走折叠的思考区tool_call渲染成工具调用卡片done触发收尾逻辑。如果没有这个字段前端只能靠猜代码会写得非常脏。2.3 分层传输层、协议层、渲染层我把整条管道分成三层各司其职这样任何一层出问题都好定位。传输层负责 HTTP 连接、SSE 帧的收发、断线重连。这一层不关心业务只保证字节流可靠地到达。协议层负责把字节流解析成StreamChunk处理粘包、半包、心跳、结束标记。渲染层在前端负责把 chunk 累积成完整消息并增量更新 DOM。分层的价值在于当出现流断了的问题时你能快速判断是传输层网络/网关的问题还是协议层解析逻辑的问题还是渲染层前端状态的问题。我见过太多项目把这三层揉在一起出问题只能靠打印日志大海捞针。3. 核心细节解析与实操要点3.1 SSE 帧格式与粘包处理SSE 的帧格式看起来简单实际暗藏玄机。一个标准的 SSE 事件长这样event: chunk data: {type:text,content:你好,index:1}注意结尾那个空行——它是事件的分隔符。服务端每发一个事件必须以\n\n结尾。问题来了TCP 是字节流不保证一次recv就拿到一个完整事件。你可能一次收到半个事件也可能一次收到三个半事件。这就是经典的粘包/半包问题。我的处理方式是维护一个缓冲区每次收到数据就追加进去然后按\n\n切分。切出来的完整事件立即处理最后一段不完整的留在缓冲区等下次数据。伪代码大概是这样buffer for raw in stream: buffer raw.decode(utf-8) while \n\n in buffer: event_str, buffer buffer.split(\n\n, 1) chunk parse_sse_event(event_str) yield chunk这里有个容易忽略的点UTF-8 多字节字符可能被切断。中文一个字占 3 个字节如果一次recv正好切在字符中间直接decode会报错。稳妥的做法是用codecs.getincrementaldecoder(utf-8)做增量解码它会自动处理跨包的字符边界。这个坑我在处理中文流式输出时踩过表现为偶尔出现乱码方块。3.2 心跳与 idle timeout 的对抗关键词里有个很扎眼的报错stream disconnected before completion: idle timeout waiting for sse。这是流式输出最经典的故障——模型在思考比如调用工具、做长推理时可能十几秒不吐任何 token中间的网关或负载均衡器一看这连接半天没数据就判定为空闲连接给掐了。解决办法是心跳。服务端在等待模型响应的间隙定期发送一个注释帧或者心跳事件: keep-alive以冒号开头的行是 SSE 的注释客户端会忽略它但它能刷新连接的活跃状态让网关认为连接还活着。心跳间隔要小于网关的 idle timeout一般设 15 秒比较稳妥因为大多数网关默认超时是 30 秒或 60 秒。注意心跳不能发太频繁否则会占用带宽、干扰前端的空闲检测。15 到 20 秒是比较舒服的区间。另外心跳帧不要带event字段避免前端把它当成业务事件处理。前端这边也要配合不能因为一段时间没收到业务 chunk 就主动断开。我的做法是前端维护一个最后收到业务数据的时间戳只有超过一个较大的阈值比如 90 秒才认为真的断了触发重连。3.3 工具调用与文本流的交织Agent 和普通聊天最大的区别是它会调用工具。这就带来一个复杂场景模型可能先输出一段文本然后发起工具调用工具执行完再继续输出文本。这些内容在流式管道里是交织的。我的处理原则是工具调用作为一个独立的事件类型不混进文本流。当模型决定调用工具时服务端发一个type: tool_call的 chunk带上工具名和参数工具执行期间服务端持续发心跳工具返回后发一个type: tool_result的 chunk然后模型继续生成发type: text的 chunk。前端收到tool_call时把当前正在累积的文本气泡封口然后渲染一个工具调用卡片收到tool_result时更新卡片状态收到后续text时开一个新的文本气泡。这样视觉上就是说话—调工具—再说话的自然节奏而不是把工具调用的 JSON 硬塞进正文里。这里有个细节tool_call的参数可能是分片到达的模型是逐 token 生成的。所以tool_call事件本身也可能有多个靠index和tool_call_id拼接。我一般让服务端在工具参数完整后再发一个tool_call_complete事件前端收到这个才开始真正执行渲染避免参数没拼完就渲染出半截 JSON。4. 实操过程与核心环节实现4.1 服务端把模型流封装成 SSE服务端的核心工作是把模型 SDK 返回的流转换成标准 SSE 帧。以 DeepSeek 的 API 为例它返回的是 OpenAI 兼容格式的流每个 chunk 长这样{choices:[{delta:{content:你},index:0}]}我要做的是把它映射成自己的StreamChunk再包成 SSE 帧。核心逻辑async def stream_agent_response(prompt, message_id): yield sse_frame(start, {message_id: message_id}) index 0 async for raw in model_client.chat_stream(prompt): delta raw[choices][0][delta] if delta.get(content): index 1 yield sse_frame(chunk, { type: text, content: delta[content], index: index, message_id: message_id, }) if delta.get(tool_calls): # 处理工具调用分片 ... yield sse_frame(done, {message_id: message_id})sse_frame是个小工具函数负责拼event:和data:行并补上结尾空行。这里我特意加了start和done两个边界事件——start让前端知道流开始了、可以准备渲染done让前端知道流正常结束了、可以收尾。没有这两个标记前端无法区分正常结束和中途断流。4.2 断流恢复从 last_index 续传流断了怎么办如果每次都从头重来用户体验很差而且浪费 token。我的方案是基于 index 的续传。前端在断流时记录下最后成功渲染的index重连时把这个index通过查询参数带给服务端。服务端如果还保留着这次生成的上下文比如把生成结果缓存了就从index1开始继续推。如果服务端已经丢了上下文那就只能重新生成但至少前端知道该从哪里清空重来。async def stream_agent_response(prompt, message_id, resume_from0): # 如果 resume_from 0 且缓存命中跳过已发送部分 ...这个机制的关键是服务端要缓存生成结果。我在 Harness 里用一个带 TTL 的内存缓存存最近几分钟的生成内容key 是message_id。这样短时间内的断流重连能无缝续上超过 TTL 就只能重来。缓存不能太大否则内存扛不住我一般限制单个会话缓存不超过 1MB。提示续传功能对移动端特别重要。手机切后台、信号抖动都会导致断流有了续传用户回来时消息能接着往下走而不是从头再来一遍。4.3 前端Vue 里的增量渲染前端这块我用 Vue 举例React 思路一样。核心是维护一个响应式的消息列表每个消息有个content字段收到 chunk 就往里追加。const messages reactive([]) function handleChunk(chunk) { let msg messages.find(m m.id chunk.message_id) if (!msg) { msg { id: chunk.message_id, content: , tools: [] } messages.push(msg) } if (chunk.type text) { msg.content chunk.content } else if (chunk.type tool_call) { msg.tools.push({ name: chunk.metadata.name, status: running }) } else if (chunk.type tool_result) { const tool msg.tools.find(t t.id chunk.metadata.tool_call_id) if (tool) tool.status done } }这里有个性能陷阱每个 token 都触发一次响应式更新会导致频繁重渲染。如果模型吐字很快一秒几十个 chunkVue 的响应式系统会被打爆页面卡顿。我的优化是批量更新——用一个缓冲区攒 chunk每 50 毫秒 flush 一次到响应式数据里。这样渲染频率从每 token 一次降到每秒 20 次肉眼完全看不出延迟但性能提升明显。let buffer let timer null function handleChunk(chunk) { buffer chunk.content if (!timer) { timer setTimeout(() { flushToMessage(buffer) buffer timer null }, 50) } }4.4 用 fetch 流式读取替代 EventSource很多人用EventSource做 SSE 客户端但它有个硬伤不支持自定义请求头。Agent 场景几乎都要带鉴权 tokenEventSource只能把 token 塞 URL 里既不安全也不优雅。所以我改用fetchReadableStream手动解析。const resp await fetch(/api/agent/stream, { method: POST, headers: { Authorization: Bearer ${token} }, body: JSON.stringify({ prompt }), }) const reader resp.body.getReader() const decoder new TextDecoder() let buffer while (true) { const { done, value } await reader.read() if (done) break buffer decoder.decode(value, { stream: true }) // 按 \n\n 切分处理 }decoder.decode(value, { stream: true })这个stream: true参数很关键它让TextDecoder内部保留跨包的多字节字符状态避免中文乱码。这个细节我在前面提过这里再强调一次因为它太容易被忽略了。用fetch还有个好处可以配合AbortController实现停止生成。用户点停止按钮时controller.abort()直接掐断请求服务端检测到连接关闭也会停止生成省 token。5. 常见问题与排查技巧实录5.1 典型故障速查表流式输出的问题五花八门我把踩过的坑整理成一张表方便对照排查。现象可能原因排查方向解决手段流中途断开报 idle timeout网关空闲超时看断流时间是否固定加心跳帧间隔小于超时中文出现乱码方块UTF-8 跨包截断检查是否用增量解码用TextDecoder的 stream 模式文本顺序错乱chunk 乱序到达打印 index 看是否递增前端按 index 缓冲排序前端卡顿每 token 触发重渲染看渲染频率批量 flush50ms 一次工具调用 JSON 半截参数分片未拼完看 tool_call 事件数等 complete 事件再渲染断流后从头开始无续传机制检查是否带 resume 参数服务端缓存 index 续传消息串台多消息流交织检查 message_id严格按 message_id 分流5.2 排查流式问题的三板斧遇到流式问题我一般按这三步走。第一步抓原始字节流。在服务端和客户端各打一份日志记录收到的原始数据。很多时候问题出在中间层网关、代理偷偷改了数据比如把\n\n换成了\r\n\r\n或者加了压缩。对比两端的日志一眼就能看出数据在哪一层被动了手脚。第二步确认边界事件。检查start和done事件是否都正常到达。如果done没到说明是断流如果start都没到说明连接根本没建立成功。这两个事件是判断流状态的锚点。第三步隔离变量。把前端换成curl直接请求接口看原始输出是否正常。如果curl正常而前端异常问题在前端解析如果curl也异常问题在服务端或中间层。这一步能快速缩小范围。5.3 几个反直觉的经验有几个经验是我踩坑之后才明白的和直觉相反但很管用。第一不要相信Content-Type。有些网关会把text/event-stream改成text/plain导致浏览器不按 SSE 处理。我的做法是前端不依赖Content-Type直接用fetch读流自己解析这样无论中间层怎么改都能工作。第二心跳不能省但也不能滥用。我见过有人为了防断流每 2 秒发一次心跳结果带宽浪费不说还干扰了前端的空闲检测逻辑。心跳的目的是骗过网关不是证明自己活着15 到 20 秒足够。第三done事件要带最终状态。不要只发一个空的done要把这次生成的完整消息 ID、token 用量、结束原因都带上。前端收到后可以更新消息状态、显示用量也方便做埋点统计。我一开始done是空的后来发现前端拿不到结束原因没法区分正常结束和被截断只能补上。第四错误也要走流。如果生成过程中出错不要直接返回 HTTP 500而是发一个type: error的 chunk然后正常关闭流。因为流已经开始了HTTP 状态码早就发出去了改不了。用 error chunk 通知前端前端能优雅地展示错误而不是白屏。6. 性能与并发让流式管道扛得住6.1 单机并发连接数的瓶颈流式输出是长连接每个活跃会话占一个连接。单机能扛多少并发取决于你的运行时。用 Python 的同步框架比如 Flask 默认模式每个连接占一个线程几百个并发就顶天了。换成异步框架FastAPI uvicorn单机扛几千个连接是常态。我在 Harness 里用的是 FastAPI 的StreamingResponse配合async生成器。关键点是生成器里不能有阻塞调用否则会卡住整个事件循环。模型 API 调用要用异步客户端数据库查询要用异步驱动任何同步的time.sleep或者同步 IO 都会拖垮并发。from fastapi.responses import StreamingResponse app.post(/api/agent/stream) async def stream(req: Request): return StreamingResponse( stream_agent_response(req.prompt), media_typetext/event-stream, )6.2 背压别让慢客户端拖垮服务端有个容易被忽略的问题如果客户端消费速度慢比如网络差而服务端生成速度快数据会在缓冲区堆积最终撑爆内存。这就是背压问题。SSE 场景下服务端一般无法直接感知客户端的消费速度但可以通过await写操作来间接实现。当底层 socket 缓冲区满了await写会挂起从而暂停生成。所以关键是用异步写不要用同步写。同步写会阻塞事件循环异步写会自然形成背压。另外服务端要设一个生成超时。如果一次生成超过比如 5 分钟还没结束强制关闭连接避免僵尸连接占资源。这个超时要比网关的 idle timeout 大否则正常的长生成会被误杀。6.3 多实例部署时的会话粘性当服务端多实例部署时续传功能会遇到麻烦断流重连可能被负载均衡打到另一个实例而那个实例没有缓存。解决办法有两个一是会话粘性让同一message_id的请求固定打到同一实例二是缓存外置把生成缓存放到 Redis 之类的共享存储里。我倾向于后者因为会话粘性在实例扩缩容时会失效。把缓存放 Rediskey 是message_idvalue 是已生成的内容和 index任何实例都能续传。代价是多一次网络往返但换来的是部署的灵活性值得。提示缓存外置要注意序列化开销。如果生成内容很大每次读写都序列化整个内容会很慢。我的做法是只缓存已发送的 chunk 列表续传时从列表里按 index 取避免重复序列化大对象。7. 我在实际项目里的一些体会流式输出这条管道写起来不难写好很难。我最大的体会是它考验的不是你对某个 API 的熟悉程度而是你对边界情况的处理能力。正常路径谁都能跑通真正拉开差距的是断流、乱序、粘包、背压这些异常场景。还有一个体会是日志要打够但要打得聪明。流式场景下日志量巨大如果每个 chunk 都打一条日志文件瞬间爆炸。我的做法是只打关键节点流开始、流结束、断流、错误、心跳超时。中间的 chunk 只在 debug 模式下打而且只打 index 和 type不打内容。这样既能定位问题又不会淹没在日志里。最后分享一个小技巧给流式管道加一个回放能力。把一次完整的流式过程所有 chunk 和它们的时间戳录下来存成文件。出问题时可以回放稳定复现。这个能力帮我定位过好几个偶发的乱序问题——线上环境难复现但回放文件一跑问题立刻现形。这个录播机制后来还被我用来做前端渲染的性能测试一举两得。