
接这个项目的时候我一开始真没把它当回事。项目内部代号 W2核心是基于 DeerFlow 智能体做二次开发对外提供对话能力而我们的团队正好需要在这个基础上封装一层 SSE 流式接口把智能体一次完整推理链路里的思考过程、工具调用、最终回答全部实时推送出去。第一版代码只用了十几行“能跑能出字”的解析逻辑前端能看到一行行文字冒出来大家都很高兴。直到测试环境开始压并发、用户端开始频繁刷新、网关开始做长连接穿透之后各种让人挠头的症状才真正冒出来。流式解析在 Demo 里永远不难真正的难度在于工程化连接怎么可靠、解析怎么容错、超时怎么判断、断线怎么恢复、日志怎么追踪、指标怎么采集。这篇文章不打算讲太多理论我直接把 W2 项目里从零做流式解析层的过程、踩过的坑、最终沉淀下来的设计完整梳理出来。适合正在做 LLM 应用、智能体二次开发或者准备把 SSE 流式接口从“能跑”提升到“能上线”的团队参考。1. 为什么基于 DeerFlow 的二次开发需要单独做一层流式封装1.1 DeerFlow 智能体输出的不是“最终答案”而是事件序列DeerFlow 这类智能体框架跟普通 LLM Chat API 最大的不同在于它的一次回答往往不是一个单体的 response。推理模型的思考、候选工具的选择、工具调用的参数序列、工具返回结果的观察再到最后的正文生成整条链路是分阶段进行的。框架内部通常会把这一过程抽象成连续的 event事件业务方如果只拿 agent 的最终返回结果相当于把智能体真正的中间状态全部扔掉了。W2 项目的诉求正好踩在这个点上我们要在聊天界面展示“正在思考”的阶段提示要在工具调用发生时看到对应的 loading 状态要让用户感知到“它不是在转圈而是在干活”。这决定了我们不能等智能体跑完再统一给数据而必须从 DeerFlow 的输出流里实时消费事件。换句话说我们在二次开发里拿到的是一个流式的上游而这个上游本身又存在多种可能的内部切换比如思考模型回答模型切换、多次工具调用、长上下文续写。这种场景下如果不把流式调用单独封装成工程层后面所有业务代码都会跟流式协议纠缠在一起。1.2 直接把流接口揉进业务代码会是什么下场我见过很多项目第一版的做法其实跟我一开始写的差不多业务 handler 里拿到一个大模型的流式响应对象然后就地循环逐行json.loads再把自己想要的那几个字段拼接起来。这套逻辑在连通性演示时完全没问题但一进测试环境就会暴露出一连串真实问题。首先解析逻辑跟业务逻辑全部耦合在一起只要网络层出现半包、乱序、空行不一致崩溃就是连锁反应。其次智能体的流式输出字段远比单个 Chat Completion 复杂内容、思考、工具调用的参数、状态变更可能都在同一个 payload 的不同字段里业务方如果只处理choices[0].delta.content要么看到思考内容被当成正文发给用户要么漏掉工具调用。再次连接异常的处理会写成无数个散落的 try-except有的人在消费端捕获超时有的人在调用端记录错误口径不一致线上日志根本没法串起来看。最后当有几个用户同时触发流式请求时资源管理会非常混乱连接没人关、事件没人消费完、任务取消也不传递最后不了了之。“流式解析工程化”本质上就是把这几件事标准化连接是连接解码是解码事件是事件业务是业务每一层的职责都清晰每一层的边界都可测。1.3 传输层、解析层、分发层的职责划分我在 W2 里最终确定了三个分层后面所有代码都是围绕这个边界展开的。层级核心职责不做的事Transport传输层建立连接、携带认证信息、配置超时、处理重连策略不关心 payload 里是什么业务字段Parser解析层把字节/文本流变成结构化事件处理 SSE 帧边界、JSON 容错不决定这个事件是“展示给用户”还是“触发逻辑”Dispatcher分发层根据事件类型做业务映射维护会话状态控制背压和广播不直接操作原始连接对象这个表写起来很简单但落实下去最容易被打破的其实是 Parser 这一层。很多人写着写着就忍不住在解析器里加业务判断比如“看到 finish_reason 就结束”或者“看到 openai 的 usage chunk 就更新 token 计数”。我的建议很明确解析器只负责产出标准化事件具体看到什么字段做什么全部放 Dispatcher 里。解析器能给出一个稳定的、可单测的事件对象后面所有工作都变得可验证。2. 流式消息解析从字节流到结构化事件的完整链路2.1 SSE 线格式与一个标准流长什么样SSEServer-Sent Events是流式接口最常用的载体。它本质上是 HTTP 响应里持续输出文本帧每一帧由若干字段行和空行构成。最常见的字段是event和dataevent不写时默认是messagedata后跟的内容就是业务负荷。帧之间用空行\n\n分隔注释行用冒号开头比如: keep-alive正常解析时应该忽略。一个典型的大模型流响应长得差不多是这样data: {id:chatcmpl-1,choices:[{delta:{role:assistant,content:},finish_reason:null}]} data: {id:chatcmpl-1,choices:[{delta:{content:你},finish_reason:null}]} data: {id:chatcmpl-1,choices:[{delta:{content:好},finish_reason:null}]} data: {id:chatcmpl-1,choices:[{delta:{},finish_reason:stop}]} data: [DONE]注意最后一行data: [DONE]它不是一个合法的 JSON 对象只是 OpenAI 兼容协议里约定俗成的结束标记。解析器的第一个必修课就是不要无脑对 data 内容做json.loads要把 “data 不是 JSON” 当成正常情况处理。2.2 智能体流式响应里常见的 delta 形态DeerFlow 二次开发场景下上游返回的 delta 通常不止 content。我实测过的响应里有很大概率出现下面几种字段choices[0].delta.content普通正文增量直接拼给回答文本。choices[0].delta.reasoning_content或thinking推理模型产出的思考增量。这玩意儿要不要给用户看是一回事但解析层必须把它单独提出来不能跟正文混在一个字符串里。choices[0].delta.tool_calls[].function.arguments工具调用的参数增量往往是分片到达的一段 JSON 字符串需要跨多个事件累积拼接后才能真正解析。choices[0].delta.tool_calls里的id和type第一次出现时经常携带这些信息后面增量只有 arguments。choices[0].finish_reason可能是stop、tool_calls、length等代表当前阶段的结束原因。如果业务侧想实现“思考之后接正文”的体验解析器必须把思考内容和正文内容作为两个独立事件发出去而不是统一叫 content。我在 W2 里定义的内部事件结构大约是这样class StreamEvent: def __init__(self, event_type, content, reasoning, tool_callsNone, finish_reasonNone): self.type event_type self.content content self.reasoning reasoning self.tool_calls tool_calls or [] self.finish_reason finish_reason这样 Dispatcher 拿到的永远是一个语义清晰的对象而不是一个裸字典。2.3 一个可以直接用的 SSE 解析器骨架下面是 W2 项目里经过多次故障打磨后的解析器核心使用 Python 的异步迭代器实现兼容主流 HTTP 客户端库import json from typing import AsyncIterator, List, Optional class ParsedSSEEvent: def __init__(self, event: Optional[str], data: str): self.event event or message self.raw_data data def json(self): try: return json.loads(self.raw_data) except json.JSONDecodeError: # [DONE]、纯文本等都走这里 return None async def parse_sse_events(stream: AsyncIterator[str]) - AsyncIterator[ParsedSSEEvent]: event_name: Optional[str] None data_lines: List[str] [] async for raw_line in stream: line raw_line.rstrip(\r) if line : # 空行意味着当前帧结束交出事件后再重置状态 if data_lines: yield ParsedSSEEvent(event_name, \n.join(data_lines)) event_name None data_lines [] continue if line.startswith(:): # 注释行一般是心跳忽略 continue if line.startswith(event:): event_name line[len(event:):].strip() elif line.startswith(data:): data_lines.append(line[len(data:):].strip()) else: # SSE 规范还允许 id、retry 等字段这里不需要就不处理 pass # 流结束但可能漏了最后一个空行仍然要把缓存中的事件交出去 if data_lines: yield ParsedSSEEvent(event_name, \n.join(data_lines))这段代码里的几个细节都是硬碰硬踩出来的用rstrip(\r)兼容\r\n行尾用空行作为帧边界多个data:行按 SSE 规范用换行拼接注释行直接跳过最后即使没有空行也要把缓存 flush 出去。调用方式也很简单配合 httpx 或 aiohttp 的文本流即可async with httpx.AsyncClient() as client: async with client.stream(POST, endpoint, jsonpayload, headersheaders) as resp: if resp.status_code ! 200: body await resp.aread() raise UpstreamError(resp.status_code, body[:512]) async for sse_event in parse_sse_events(resp.aiter_lines()): # 这里只拿到事件不做业务判断 handle_event(sse_event)2.4 跨 chunk 字符、坏 JSON 与不标准帧的容错真正做工程化时最让人头疼的不是“标准 SSE”而是各种笨拙的上游实现和网络路径。比如有些网关会把一个 SSE 帧的data:内容拆成多个 TCP chunk 送过来客户端如果按包处理内容就可能把 JSON 截断成两半。解决方法是始终用 HTTP 客户端的aiter_lines或aiter_text拿到已经拼接好的文本行而不是自己直接处理aiter_bytes。因为aiter_text在底层使用了增量解码器天然处理了 UTF-8 字符跨越多个 chunk 的边界问题。如果你非要自己处理裸字节记得用codecs.getincrementaldecoder(utf-8)()做增量解码否则中文内容只要恰好被拆在 chunk 边界上就是乱码。另一个常见问题是“数据行里包含的不是合法 JSON”。data: [DONE]是最典型的合法非 JSON 场景此外某些平台还会在流里插入状态文本比如data: connected。解析器对这类内容绝不能直接抛异常应该原样保存在ParsedSSEEvent.raw_data里由上层判断。还有更离谱的情况上游把两条事件挤在同一个 chunk 里中间只有换行没有空行。遇到这种实现纯按空行切帧就会把两条事件粘在一起。稳妥的做法是在解析器外层再做一层“尝试从拼接数据里提取完整 JSON 对象”的兜底逻辑。这里可以复用 JSONDecoder 的raw_decode方法持续提取import json def extract_json_objects(text: str): decoder json.JSONDecoder() idx 0 length len(text) while idx length: while idx length and text[idx] not in {[: idx 1 if idx length: break try: obj, end decoder.raw_decode(text[idx:]) yield obj idx end except json.JSONDecodeError: break这个函数只能作为兜底不推荐当主解析路径。最理想的情况还是上游按规范输出但工程上必须做好“上游不规范”的准备。提示解析层的核心指标是“任何输入都不能让进程崩溃”。哪怕遇到完全不可解析的二进制垃圾也应该以错误事件形式抛给上层而不是在解析器里直接抛异常终止整个任务。3. 封装 SSE 流式接口调用逻辑传输层与生命周期的打磨3.1 代码结构上的分层设计W2 项目里最终的调用链路大概是这样的业务请求入口 → AgentStreamSession对外暴露 async for 事件的迭代器 → HttpTransport创建连接、设置认证、管理超时 → SSE Parser把文本流转成 ParsedSSEEvent → EventDispatcher把事件按类型分发维护 tool_call 聚合 → 消费端WebSocket 推送、消息队列投递等每一层只依赖下一层提供的抽象接口。比如 Parser 不关心是 httpx 还是 aiohttp 产生的流只接收一个异步迭代器。HttpTransport 不关心解析出的内容是正文还是工具调用只负责把连接管理好。这样后续即使把 DeerFlow 的上游从 HTTP 换成内部 gRPC也只需要换 Transport解析和分发基本不动。3.2 一个请求就是一个异步生成器封装 SSE 调用逻辑时最自然的抽象是“异步生成器”。因为生成器天生支持async for也天然表达“事件是逐条到达的”。更关键的是异步生成器在aclose()或被垃圾回收时会自动触发finally块保证底层连接被释放。W2 核心的调用骨架大概是这样的async def create_chat_stream(session: AgentSession, payload: dict): config session.config headers { Authorization: fBearer {config.api_key}, Accept: text/event-stream, } timeout httpx.Timeout( connectconfig.connect_timeout, readconfig.read_timeout, writeconfig.write_timeout, poolconfig.pool_timeout, ) async with httpx.AsyncClient(timeouttimeout) as client: async with client.stream(POST, config.endpoint, jsonpayload, headersheaders) as resp: if resp.status_code ! 200: body await resp.aread() raise UpstreamError( fupstream returned {resp.status_code}, body{body[:512]!r} ) async for parsed in parse_sse_events(resp.aiter_lines()): if parsed.raw_data [DONE]: return yield parsed这个代码看似简单但里面有三处必须坚定一是async with client.stream保证无论如何退出都会关闭连接二是[DONE]在解析器里不会导致 JSON 异常在这里做结束判断也顺理成章三是网络异常在生成器内部自然上抛由外层统一处理。3.3 生命周期状态机与取消传播一个完整的流式调用生命周期至少应该包含这些状态IDLE、CONNECTING、STREAMING、DONE、ERROR、CANCELLED。我在工程里封装了一个StreamSession对象它对外暴露state和事件流两个东西。状态机的作用不是摆设而是让上层能准确判断“当前是连接还没建立、正在流式输出、还是已经异常结束”。class StreamState(Enum): IDLE idle CONNECTING connecting STREAMING streaming DONE done ERROR error CANCELLED cancelled取消传播是被严重低估的点。当客户端断开 WebSocket 时应用层应该主动取消本次生成器。如果直接对生成器调用aclose()Python 会在生成器当前挂起的位置注入GeneratorExit或asyncio.CancelledError。如果你在业务代码里捕获了这个异常又没有重新抛出取消就静默失效了底层连接还会继续被拉取白白浪费上游资源。这一点我在后面的故障复盘里还会提到它是很少见但极其致命的 bug。3.4 超时、心跳与空闲慢思考的处理流式接口的超时配置跟普通 HTTP 请求完全不是一个思路。普通请求讲究总超时流式请求要是也弄个“30 秒总超时”只要模型思考超过 30 秒就会在中途被掐断。W2 里的配置如下connect_timeout3 秒建连超过这个时间直接失败。write_timeout10 秒请求体发送超时。read_timeout120 秒这是“两次 read 之间没有数据到达”的超时不是总超时。这个read_timeout是安全阀如果模型思考时间过长一直不吐 token连接就会一直安静客户端不能无限等下去。但实际体验又不能 120 秒没有任何反馈所以需要在 Dispatcher 层做“心跳事件”如果距离上一次收到上游数据已经超过 10 秒就往消费端推一个“waiting”状态前端看到的效果就是“它还在思考中”。这个心跳事件绝不能由 Parser 层生成因为 Parser 只认上游字节流不关心用户体验。对于长思考场景更精细的做法是分两级首字节超时TTFT用 20 秒超过后告警首字节之后的闲置读取超时用 120 秒。这种区分看着很小但能帮你在线上快速区分“服务卡死”和“单纯响应慢”。3.5 重试到底该不该做以及怎么重试流式接口的重试是一个大坑。很多人直接把普通 HTTP 的重试策略套过来结果把简单问题搞复杂。问题在于流式请求不天然幂等如果上游已经把思考内容和工具调用发出去了你这边因为网络抖动发起重试新连接会重新生成一遍同样的内容。工具调用场景下更危险重试可能会触发两次真实的工具副作用。我在 W2 里定的原则是还没有收到任何事件数据也就是首字节之前发生错误可以安全重试重试次数一般 2~3 次用指数退避。已经收到至少一个事件后再断开绝不自动重试。此时应该把连接中断作为StreamInterrupted错误抛给上层由上层决定是让用户手动重试还是展示“断线重连”按钮。如果要在上层做自动恢复就自己去重整个会话上下文向新连接提交一个“从第 N 个事件开始续传”的请求。但绝大多数上游不提供这种能力所以这个方案通常只能做软恢复比如前端重新进入“loading”状态。提示给每个流式请求配一个stream_id重试时使用不同的 ID。否则日志/追踪系统里多条流的 ID 完全相同线上根本没法区分哪次是重试前的残流、哪次是重试后的新流。4. 流式解析工程化里避不开的可观测性、并发与测试4.1 用 stream_id 把一条流的生命周期串起来流式请求往往跨越多个服务浏览器触发、接入层转发、应用服务调用 DeerFlow、DeerFlow 再调底层模型。任何一个环节出问题如果日志无法关联排查就直接从地狱模式开始。W2 里我给每条流分配了一个stream_id格式类似stream_hex16从应用服务入口生成塞进 HTTP header、日志字段、追踪 span、指标标签里。日志只靠观察是看不出来的我通常强制要求每一行日志都带 stream_id[stream_2fa1c9] stateconnecting url/v1/agent/chat attempt1 [stream_2fa1c9] statestreaming seq5 delta_bytes96 reasoning_bytes128 tool_calls0 [stream_2fa1c9] statedone reasonstop total_events87 ttft_ms412 duration_ms18223这套日志格式的好处是用 grep 一搜就能把一条流从生到死的全链路拉出来。流式接口的排错速度靠的就是这个关联能力。4.2 指标别只盯 QPS这些才是流式服务的核心指标流式服务的健康度跟普通 HTTP 完全不同。普通服务看错误率、延迟分位流式服务还要看下面这些指标含义作用ttft_ms从发起到收到第一个事件的时间反映“智能体开始说话”的速度stream_duration_ms整条流从开始到 DONE 的耗时反映一次完整回答的体验events_total一条流里的事件数判断是否存在事件风暴interrupted_total上游关闭但未收到 DONE 的次数最重要的稳定性预警指标parse_error_total解析器遇到非法内容后转为错误事件的次数发现上游格式变更reconnect_total首字节前重试次数反映网络/网关稳定性其中interrupted_total是我最看重的因为“无声中断”是流式系统独有的故障模式普通错误率根本捕捉不到。用户看到答案只生成了一半但在服务端没有 4xx/5xx也没有异常堆栈只有这个计数器在涨。4.3 消费端背压与并发控制流式输出的生产者速度通常远快于消费者能处理的速度特别是在业务侧还要做工具调用聚合、WebSocket 推送、落库记录的时候。如果不做背压控制内存里的事件堆积会把进程拖垮。我在 W2 里用的方案是asyncio.Queue。解析器拿到事件后往队列写消费端从队列读取。队列设置了maxsize默认 1000。队列满时解析器协程阻塞在put上间接把背压传导给上游流。如果消费端明确已断开就调用queue.put_nowait(StreamClosedSentinel())并取消生成器避免继续消费上游流量。DeerFlow 智能体会话本身可能是有状态的同一会话最好不要跨并发请求混用。我建议给每个会话建独立的Semaphore(2)或Semaphore(1)并发请求超限就直接返回“会话忙”而不是并发打爆上游。这一步纯属资源治理却决定了系统能不能扛住多人同时调用。4.4 用测试把协议行为钉死流式解析的回归测试是必须的。你不可能每次都拿真模型来测试解析器那样既慢又贵。正确做法是用录制好的流式响应片段作为 fixtures喂给解析器和 Dispatcher。一个最小测试长这样async def test_parse_multi_data_lines(): payload [ data: {choices:[{delta:{content:你}}]}\n, \n, data: {choices:[{delta:{content:好}}]}\n, \n, data: [DONE]\n\n, ] events [] async for event in parse_sse_events(iter(payload)): events.append(event) assert len(events) 3 assert events[0].json()[choices][0][delta][content] 你除了这种 happy path测试矩阵里一定要包含思考字段与正文混出、工具调用 arguments 跨多个事件、[DONE]前没有空行、非法 JSON 文本、流中途断开的 EOF、带心跳注释行的长连接。契约测试也很重要。把上游每种事件格式的 JSON 样本保存到仓库每当 DeerFlow 框架升级或换了模型供应商就跑一次合约测试任何字段结构变化都能第一时间暴露。说实话这个测试帮我避过好几次上游悄悄改字段的坑。4.5 接入层长连接配置的几个常识SSE 服务基本都会穿越接入层到达外部这里有几个让无数团队翻车的配置项。第一响应缓冲必须关闭。接入层默认会把上游响应攒够一定大小再发给客户端这对 SSE 是灾难会造成用户端长时间收不到一个字节表现就是“连接还活着但就是不出字”。关闭响应缓冲让数据以 chunked transfer 方式实时透传。第二接入层的读取超时不能设置得太小。长思考场景下整条流可能持续几分钟如果接入层在 60 秒无新数据就掐断下面服务再稳定也白搭。建议至少调到 300 秒或者干脆关闭该路径的读超时仅依赖应用层的闲置超时来兜底。第三不缓存流式响应。CDN、缓存中间件如果开了缓存会把整条流缓冲完再统一释放这让用户等到怀疑人生。SSE 路径必须在缓存规则里显式跳过。5. 一场线上“半个回答”事故的完整排查复盘5.1 事故现象连接还悬着回复却像被腰斩某个周末线上突然有用户反馈机器人回答到一半就停了前端一直显示“正在输入”但没有任何新内容到达。重启服务没解决用户刷新会话重新提问后偶尔又正常属于那种最烦人的“偶发必现却又无法随时复现”的故障。我第一时间去看了应用服务的错误日志。奇怪的是服务端一条 ERROR 都没有。既没有超时记录也没有连接被拒绝的异常日志里最后一条状态还是streaming。这说明整个生成器仍在正常运行事件还在被消费但下游用户没有收到。最可怕的就是这种“看起来一切正常实际上已经断了”的状态。5.2 排查链路从指标计数到抓包逐层收窄我先查了interrupted_total和stream_duration指标。异常流的stream_duration都集中在一个非常接近的数值附近大概 60 秒。这个数字太整齐了不像是随机网络抖动更像是某个固定的超时阈值。再往前推应用服务到外网用户之间还有一层接入层网关默认读超时恰好就是 60 秒。到这里嫌疑已经锁定了大半。接着我拉取了异常流的追踪链路时间线非常清晰连接建立后约 2 秒上游正常返回第一个事件持续收到约 50 秒的递增内容之后整条流进入了一个“模型长思考”的静默段没有产生任何 token也没有心跳第 60 秒接入层因为没读到新字节主动掐断了连接。掐断后应用服务并没有立刻感知到异常因为底层 TCP 连接是优雅关闭的生成器继续挂在那里等待新事件直到等待超时才会报错。用户侧看到的就是“不结束、不报错、也不出字”。我用 tcpdump 在应用节点和网关之间抓包验证了一下断连方向确实是接入层发的 FIN时间也精确落在 60.0x 秒附近。至此根因基本确定接入层读超时 长时间无字节 服务端缺失心跳保活机制三个条件共同制造了这起事故。5.3 根因修复三个动作缺一不可修复分三层同时进行。第一层接入层调整。把流式路径的读超时从 60 秒调到 600 秒并显式关闭响应缓冲。这个改动从根本上解除了“长思考即被误杀”的定时炸弹。第二层应用服务增加保活机制。在 Dispatcher 层如果发现距离上一次收到上游事件已经超过 10 秒就主动生成一个: keep-alive注释帧或者一段waiting事件推到下游。注释帧的代价几乎为零但能保证连接每 10 秒至少有一个字节经过网关彻底规避“静默被杀”。第三层也是我认为最重要的一层修改 StreamSession 的错误语义。以前上游关闭连接但服务端没收到[DONE]时生成器的表现跟正常流结束一样只是不再有新事件。我改成在async for循环自然退出前检查是否收到了DONE标记没收到且状态不是正常结束时统一抛出StreamInterrupted。Dispatcher 捕获后一方面把interrupted_total指标加一另一方面向前端发送一个明确的“流已中断”事件前端可以据此展示“重新连接”按钮而不是永远转圈。5.4 如何防止类似问题重演这次事故之后我补了两类检查。一类是持续压测。我写了一个“长静默复现脚本”控制 mock 上游在输出 50 秒内容后静默 70 秒不发任何字节然后观察整条链路是否会触发StreamInterrupted。只要这个脚本在 CI 里每天跑一遍任何新加入的接入层超时或代码退化都能第一时间暴露。另一类是告警规则。对interrupted_total设置增量告警阈值是每 5 分钟大于 0 就触发 P2。这个指标比任何“接口错误率”都更能反映流式系统的真实健康度因为它捕捉的是“连接还在但内容断了”的隐性故障。6. 沉淀下来的几条反直觉经验6.1 流式重试最好的时机是“第一条事件还没到”的时候传统接口把“重试”做成遇到错误就立刻回调但在流式世界里这几乎一定制造重复。原因是工具调用、思考过程这些增量事件不是幂等的你跑到一半再重连等于让智能体从头跑一遍用户不仅会看到重复回答还可能因为重复的工具调用触发真实副作用。所以我把重试逻辑严格限定在“首字节之前”发生的错误上。一旦有过任何事件发出就只做中断上报不再自动重试。6.2 解析器应该吐“原样”别在解析层做业务拼接写解析器的人很容易顺手把content和reasoning_content拼成一个字段再传给上层。这个动作看似方便实际是把业务决策提前到了协议层。思考内容要不要给用户看、什么时候给用户看是 Dispatcher 甚至前端才该决定的事。解析器一旦开始“好心”拼接后续任何人想区分这两段内容都必须重新解析原始数据而原始数据又未必保留了。保持原样输出不丢失信息这是流式数据管道最重要的原则。6.3 日志和原始帧的保留是线上排障的最后底牌流式接口的 bug 往往是概率性的不好复现。如果日志只记录“处理后的状态”原始事件长什么样就无从知晓了。W2 里我专门加了一个 debug-trace 模式开启后会把每条流收到的原始 SSE 行脱敏后写入独立的 trace 文件。平时不开启只有遇到疑难杂症时对指定 stream_id 开启。这个模式帮我定位过上游返回非法 JSON、解析器被意外空行截断、内容字段被莫名加了前缀等至少三个隐蔽问题。6.4 不要在 finally 块里吞掉 CancelledError最后一个我要写的经验可能也是所有异步系统里代价最隐蔽的一个不要在清理逻辑里吞掉CancelledError。很多人写生成器时会写except Exception包住清理逻辑但asyncio.CancelledError在 Python 3.8 之后继承自BaseException它根本不会被except Exception捕获。最怕的是另一种写法用except asyncio.CancelledError: pass吃掉取消清理完还继续执行后面的代码。这会让取消完全失效底层连接不断被拉取上游智能体白白产生事件流直到超时或连接池耗尽才暴露。正确做法是捕获取消后完成必要的清理然后raise重新抛出。让取消信号继续向调用链上层传播该关的连接关掉该断的会话断掉所有状态状态都收敛到 CANCELLED。这个细节不写下来真的很容易忘而它往往是并发一上来才出现的“神秘卡死”的元凶。说实话做完 W2 这个项目以后我对“流式解析工程化”的理解已经完全变了。它不再是一段解析data:行的代码而是传输层、解析层、分发层、可观测性、测试、资源治理串起来的一整条链路。如果你现在正准备接一个智能体流式接口我的建议很简单先按这个分层把骨架搭好再补上中断计数和原始帧日志最后再开始写业务界面。顺序反了后面大概率会回头重构一遍而且是在线上故障的重压下重构。