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通过查询参数带给服务端。服务端如果还保留着这次生成的上下文(比如把生成结果缓存了),就从index+1开始继续推。如果服务端已经丢了上下文,那就只能重新生成,但至少前端知道该从哪里清空重来。
async def stream_agent_response(prompt, message_id, resume_from=0): # 如果 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 都触发一次响应式更新,会导致频繁重渲染。如果模型吐字很快,一秒几十个 chunk,Vue 的响应式系统会被打爆,页面卡顿。我的优化是批量更新——用一个缓冲区攒 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 场景几乎都要带鉴权 token,EventSource只能把 token 塞 URL 里,既不安全也不优雅。所以我改用fetch+ReadableStream手动解析。
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 触发重渲染 | 看渲染频率 | 批量 flush,50ms 一次 |
| 工具调用 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_type="text/event-stream", )6.2 背压:别让慢客户端拖垮服务端
有个容易被忽略的问题:如果客户端消费速度慢(比如网络差),而服务端生成速度快,数据会在缓冲区堆积,最终撑爆内存。这就是背压问题。
SSE 场景下,服务端一般无法直接感知客户端的消费速度,但可以通过await写操作来间接实现。当底层 socket 缓冲区满了,await写会挂起,从而暂停生成。所以关键是用异步写,不要用同步写。同步写会阻塞事件循环,异步写会自然形成背压。
另外,服务端要设一个生成超时。如果一次生成超过比如 5 分钟还没结束,强制关闭连接,避免僵尸连接占资源。这个超时要比网关的 idle timeout 大,否则正常的长生成会被误杀。
6.3 多实例部署时的会话粘性
当服务端多实例部署时,续传功能会遇到麻烦:断流重连可能被负载均衡打到另一个实例,而那个实例没有缓存。解决办法有两个:一是会话粘性,让同一message_id的请求固定打到同一实例;二是缓存外置,把生成缓存放到 Redis 之类的共享存储里。
我倾向于后者,因为会话粘性在实例扩缩容时会失效。把缓存放 Redis,key 是message_id,value 是已生成的内容和 index,任何实例都能续传。代价是多一次网络往返,但换来的是部署的灵活性,值得。
提示:缓存外置要注意序列化开销。如果生成内容很大,每次读写都序列化整个内容会很慢。我的做法是只缓存"已发送的 chunk 列表",续传时从列表里按 index 取,避免重复序列化大对象。
7. 我在实际项目里的一些体会
流式输出这条管道,写起来不难,写好很难。我最大的体会是:它考验的不是你对某个 API 的熟悉程度,而是你对边界情况的处理能力。正常路径谁都能跑通,真正拉开差距的是断流、乱序、粘包、背压这些异常场景。
还有一个体会是,日志要打够,但要打得聪明。流式场景下日志量巨大,如果每个 chunk 都打一条,日志文件瞬间爆炸。我的做法是只打关键节点:流开始、流结束、断流、错误、心跳超时。中间的 chunk 只在 debug 模式下打,而且只打 index 和 type,不打内容。这样既能定位问题,又不会淹没在日志里。
最后分享一个小技巧:给流式管道加一个"回放"能力。把一次完整的流式过程(所有 chunk 和它们的时间戳)录下来,存成文件。出问题时可以回放,稳定复现。这个能力帮我定位过好几个偶发的乱序问题——线上环境难复现,但回放文件一跑,问题立刻现形。这个录播机制后来还被我用来做前端渲染的性能测试,一举两得。