NestJS + LangChain + SSE:AI Agent流式事件推送最佳实践

NestJS + LangChain + SSE:AI Agent流式事件推送最佳实践 AI Agent 开发落地的关键一环NestJS LangChain 把流式事件推给前端的正确姿势去年年底我接了一个内部知识库问答 Agent 的项目后端用的 NestJSAgent 编排层用的 LangChain。功能本身不复杂用户提问Agent 调用检索工具把答案流式地吐给浏览器。可真正动手时才发现最卡脖子的根本不是 Agent 的推理链路怎么搭而是“后端算出来一堆事件怎么稳定地、实时地送到前端”。HTTP 一把梭返回全文用户等十秒看一个转圈体验直接归零用 WebSocket 又觉得重毕竟我们根本不需要客户端往服务端推东西。最后落地方案是 SSE也就是 Server-Sent Events配合 NestJS 的Sse()端点和 LangChain 的流式迭代器把 Agent 的思考过程、工具调用、token 增量一条条地推过去。这篇文章不聊 PPT 架构图就聊聊这台链路是怎么一步步搭起来的中间踩了哪些坑以及从 Demo 到生产环境还需要处理哪些要命的小问题。适合正在用 NestJS 做 AI 应用后端、或者想把 LangChain Agent 输出做流式化的同学参考。1. 为什么是 SSEAI Agent 场景下最好的信息通道不是 WebSocket1.1 传统请求/响应模式在 Agent 场景下的三个翻车现场先说什么场景下你会意识到必须上流式。我第一个版本的接口很朴素前端 POST 一个问题后端跑 LangChain Agent等到整个 answer 拼接完再一次性 return。三个问题接踵而至一个 Agent 调用链可能包含 LLM 推理、工具检索、再推理总耗时动辄 5~15 秒。Nginx 默认proxy_read_timeout是 60 秒本地没问题部署到服务器上用户稍微多问几轮连接就被网关掐了前端拿到 504。用户体验非常差。用户看到的是一个静止的“发送中”状态他不知道 Agent 是卡死了还是在思考只能干等。几乎每个试用的人都会忍不住再点一次发送造成重复请求Agent 被同一个问题打两遍。调试痛苦。LangChain 在推理过程中会产生大量中间事件比如调用了哪个工具、检索到了几条文档、LLM 输出了什么中间 token。这些信息在普通响应模式里要么丢掉要么得靠日志事后翻没法在界面上实时展示。所以流式不是“锦上添花”是做 AI Agent 应用的硬需求。你至少得让用户看到“草在动”他才知道马在吃草。1.2 SSE 和 WebSocket 的取舍单向推送就该用 HTTP 流很多人一听到“实时推送”就条件反射上 WebSocket但 AI Agent 这个场景里SSE 明显更合适。WebSocket 是一个全双工的长连接协议双方可以随时互发消息。能力确实强可也带来了几个麻烦需要单独维护连接状态、需要处理心跳和重连逻辑、在浏览器端需要写onopen/onmessage/onclose一堆回调在 Nginx 层还得专门配置 Upgrade 头。而 SSE 的本质是 HTTP 长连接服务端往客户端单向推天然契合 Agent 的输出场景客户端发一个问题服务端连续推一串事件最后推一个结束事件把连接关上。它有几个点名表扬的特性基于普通 HTTP不需要额外协议NestJS 的 Controller 加一个装饰器就能用Nginx 几乎零配置。自带Last-Event-ID断线重连机制浏览器EventSource对象内置重连能力前端代码量少得可怜。事件格式是标准化的text/event-stream用event:和data:字段区分事件类型和内容调试时直接 curl 就能看数据。你在 Agent 场景下根本不需要“前端推给后端”这个反向通道真要传中止信号发一个普通 POST 请求就行。用 WebSocket 属于杀鸡用牛刀还把运维复杂度拉高了。1.3 理解 SSE 的数据格式和 NestJS 的最小实现SSE 的协议本身简单到让你怀疑人生本质就是一段持续输出的文本流event: message data: {content:你好} event: done data: {id:conv_123}每两个\n\n之间就是一个事件。event字段是事件名data字段是内容可以是任意字符串我们约定用 JSON。连接建立后服务端返回Content-Type: text/event-stream然后持续往响应里写这段文本即可。NestJS 对 SSE 的支持非常成熟Controller 里直接写Sse(agent/stream) streamAgent() { return new Observable((subscriber) { subscriber.next({ data: { content: hello } }); subscriber.complete(); }); }Sse()装饰器会自动设置好响应头返回值是一个 RxJSObservableMessageEventMessageEvent的形状就是{ data: any, id?: string, event?: string, retry?: number }。你不需要手动去管req和res框架替你处理了。这个设计非常贴合 Node.js 生态因为 LangChain 的流式输出本质上是异步迭代器而 RxJS 的 Observable 天然能把异步迭代器转成事件流。后面我会给具体接法。2. 先把 Agent 的流式链路跑通LangChain 端的事件产出逻辑2.1 当前版本的技术选型LangChain 0.3.x 与 LangGraph 的取舍聊 NestJS 之前得先把上游的 Agent 链路捋顺。我用的 LangChain 版本是 0.3.x 系列。这个版本有一个重要变化老的initializeAgentExecutorWithOptions、AgentExecutor那套 API 已经逐渐退居二线官方更推荐用 LangGraph 来构建 Agent或者用create_react_agent、create_tool_calling_agent这类工厂函数。也不是说老 API 不能用了只是如果你的项目是全新的建议直接走新链路。LangGraph 的核心优势在于它把 Agent 的节点和边显式化了模型节点、工具节点、条件边你能看到每一步的执行情况也能方便的挂载各种中间回调。这个特性跟 SSE 是绝配——节点切换时你可以发一个事件出去用户一眼就能看到 Agent 正在做什么。我这里给一个基于create_react_agent的示例它内部其实就是 LangGraph 的封装适合大多数“模型 工具”的 Agent 场景。2.2 让 LLM 的 token 增量真正流出来LangChain 的模型层对流式支持是默认的。用langchain/openai这个包时你只要在调用链里传入streaming: true或者直接调用stream()方法就能拿到一个AsyncGenerator。关键代码长这样import { ChatOpenAI } from langchain/openai; import { createReactAgent } from langchain/langgraph/prebuilt; import { Tool } from langchain/core/tools; const model new ChatOpenAI({ model: gpt-4o-mini, temperature: 0.2, streaming: true, }); const searchTool new Tool({ name: knowledge_search, description: 在内部知识库中搜索相关内容, func: async (query: string) { // 实际会调用向量检索服务 return 检索结果...; }, }); const agent createReactAgent({ llm: model, tools: [searchTool], });有了agent之后调用它的stream()方法你会拿到一长串 LangGraph 的事件对象。这些事件对象有统一的类型标记比如“进入节点”“离开节点”“模型产出了新 token”“工具返回了结果”。我强烈建议你先写一段脚本把它打出来看看眼见为实而不是凭空猜。2.3 事件归一化LangGraph 事件与前端展示的桥接LangGraph 的原始事件结构比较复杂直接丢给前端显然不行。我会在 Agent 层做一个事件归一化层把原始事件转成三类统一结构agent:startAgent 开始执行。agent:messageLLM 产出的增量 token前端拿它逐字拼接。agent:tool_callAgent 决定调用某个工具前端可展示“正在检索...”。agent:tool_result工具返回结果前端可展示“检索到 N 条内容”。agent:done整个推理结束附带最终 answer。agent:error链路异常。归一化的代码我习惯放在一个独立的 service 里比如AgentStreamService。它的核心是一个循环不断从 LangGraph 的流里取出事件映射成统一结构再通过一个Subject发射出去Injectable() export class AgentStreamService { async streamAgent(query: string): PromiseAsyncGeneratorAgentStreamEvent { const stream await agent.stream( { messages: [{ role: user, content: query }] }, { streamMode: messages } ); async function* generate() { for await (const [messageChunk, metadata] of stream) { if (typeof messageChunk.content string) { yield { event: agent:message, data: { content: messageChunk.content }, }; } } } return generate(); } }这里用streamMode: messages拿到的是模型 token 级别的产出如果你想拿节点级别的执行步骤就用默认的streamMode: updates。两者可以组合不过最初实现建议先只跑通 messages 模式后面再加节点事件。有一点要注意agent.stream()返回的是一个异步生成器它天然适合被消费者消费。设计 service 接口时我倾向于直接返回AsyncGeneratorAgentStreamEvent而不是把 Observable 的语义提前引入业务层——异步生成器在 NestJS 里可以直接转成 SSE 流后面你会看到这个转换有多顺滑。2.4 工具调用的流式事件怎么处理如果 Agent 只是单纯调一次模型那这个流式链路其实很简单。但真实 Agent 一定会调用工具。工具的执行时间是不确定的可能是 200ms 的 API 调用也可能是几十秒的数据分析任务。这段时间里用户如果看不到任何反馈他大概率会觉得服务挂了。所以我会在“模型决定调用工具”和“工具执行完成”这两个节点分别发一个事件。LangGraph 里怎么做最直接的办法是在工具执行前后包一层function withStreamingTool(tool: Tool) { return new Tool({ name: tool.name, description: tool.description, func: async (input, config) { // 此处可通知前端工具调用开始 const result await tool.invoke(input, config); // 此处可通知前端工具调用完成 return result; }, }); }不过更干净的做法是直接监听 LangGraph 的节点事件通过config里传入的callbacks捕获on_tool_start和on_tool_end。但这个回调比较底层封装起来略繁琐。我实际项目里更常用的方式是在createReactAgent调用后手动维护一个“工具调用状态”变量结合 LangGraph 的事件流判断当前停在哪一步。如果你的 Agent 相对复杂建议直接用 LangGraph 的StateGraph显式定义一次工具调用节点这样不但能在节点入口发事件还能做更精细的异常处理比在工具函数里塞回调要清爽得多。3. NestJS 的 SSE 端点把异步生成器变成浏览器能读的流3.1 Controller 层设计别在请求线程里跑 Agent服务端代码要拆成两层Controller 只负责建立 SSE 通道业务逻辑全部委托给 service。我踩过的一个典型坑是在 Controller 的方法体里直接await agent.stream()导致整个 Handler 被阻塞。在 SSE 场景下这个“阻塞”意味着浏览器收到的是一段空白流直到所有事件攒齐才一次性输出流式效果完全丢失。正确做法是Controller 方法不await业务逻辑而是返回一个 RxJSObservable由 NestJS 在订阅时再去执行耗时逻辑。这样 HTTP 响应头立刻返回连接建立后事件从 Observable 里逐个发射出去。3.2 用 RxJS 桥接异步生成器的两种写法写Observable的时候你可以选择用rxjs的from转异步生成器也可以手写 subscriber。我推荐先看手写版本因为它的执行时机更可控import { Controller, Sse, MessageEvent } from nestjs/common; import { Observable } from rxjs; import { AgentStreamService } from ./agent-stream.service; Controller(agent) export class AgentController { constructor(private readonly agentStream: AgentStreamService) {} Sse(stream) stream(Query(question) question: string): ObservableMessageEvent { return new ObservableMessageEvent((subscriber) { const run async () { try { subscriber.next({ event: agent:start, data: { ts: Date.now() }, }); const generator await this.agentStream.streamAgent(question); for await (const chunk of generator) { subscriber.next({ event: chunk.event, data: chunk.data, }); } subscriber.next({ event: agent:done, data: { ts: Date.now() }, }); subscriber.complete(); } catch (err) { subscriber.error(err); } }; run(); }); } }这段代码的核心在于new Observable((subscriber) { ... })是惰性的只有当 NestJS 内部真正订阅它时里面的run()才会执行。这个订阅时机恰好是 HTTP 响应头已经发出、SSE 连接建立完毕的时候所以浏览器会先收到agent:start接着就是一连串的 token 增量事件。如果你更喜欢声明式写法可以这样import { from } from rxjs; import { mergeMap } from rxjs/operators; return from(this.agentStream.streamAgent(question)).pipe( mergeMap((generator) from(generator)), map((chunk) ({ event: chunk.event, data: chunk.data })), );两种写法本质等价但手写版本方便你在中间插入“心跳”和“错误兜底”。我线上用的是手写版本因为长期连接最容易出的问题不是逻辑错而是中间某个环节卡住没有新数据也没有报错连接就僵在那里。手写版本里你可以塞一个区间定时器每 15 秒subscriber.next({ event: heartbeat, data: { ts: Date.now() } })确保连接保持活性。3.3 心跳、注释行和连接关闭的边界情况SSE 协议有一个经常被忽略的细节如果服务端长时间不发送任何数据包括空行某些代理服务器或浏览器会判定连接超时。这就是热搜词里 “before completion: idle timeout waiting for sse” 这个报错的典型来源之一。规避方式有两层应用层每 15 秒发送一个heartbeat事件。协议层每隔几秒输出一个: keep-alive\n\n注释行。SSE 规范里以冒号开头的行是注释浏览器会忽略但能起到刷新 TCP 连接活动状态的作用。NestJS 里实现心跳直接在 Observable 里维护一个setInterval并在subscriber.complete()或subscriber.error()时clearInterval。注意一定要清理否则连接关了定时器还在跑会导致内存泄漏。连接关闭的另一个边界情况是客户端主动断开了。浏览器关页面、切换路由、断网连接就断了。如果你只闷头往下游 LLM 发请求会浪费大量 token。这个问题我会在第 5 节单独展开。4. Agent 事件协议设计前端收到的是什么以及怎么解析4.1 一套够用的事件类型定义后端和前端之间必须有一套约定。我见过很多项目在这个地方偷懒直接把 LangChain 的原始事件发给前端导致前端代码耦合了大量框架概念。我的建议是定义一套稳定的应用层协议LangChain 内部怎么变化都不影响前端。最终事件类型如下event 字段data 结构含义agent:start{ ts: number }链路开始agent:message{ content: string }LLM 增量 tokenagent:tool_call{ name: string, input: any }工具开始调用agent:tool_result{ name: string, output: string, duration: number }工具返回agent:done{ ts: number }链路结束agent:error{ message: string }异常信息heartbeat{ ts: number }心跳保活这套协议的设计原则是粒度适中。不要细到把 LangGraph 的每个节点状态都透传出去也不要粗到只有一段完整答案。对前端而言中间过程展示是有价值的但展示的层级最好控制在“思维链摘要 工具状态”而不是底层 executor 的每一步内部流转。4.2 SSE 中的 JSON 转义和格式坑SSE 协议要求data:后面跟的是一行文本但其实可以是多行。如果你直接塞 JSONJSON 里又带了换行符会导致事件被切割成两个事件前端解析时直接报错。解决方式是发送前对 JSON 做单行化。比如把JSON.stringify(chunk.data)里的\n替换成\\n或者用\u000A之类的转义。这个坑不测一次你是不会注意到的尤其是 LLM 输出的 markdown 文本天然包含大量换行我第一次联调时前端收到的内容就是断的。NestJS 的MessageEvent对象内部其实会帮你处理一部分序列化只要你的data字段是字符串它默认会放在data:后面。但多行问题它不管所以最稳妥的办法是在 service 层统一把 data 转成单行字符串再传给 Observable。4.3 前端解析用 EventSource 还是 fetch ReadableStream浏览器原生支持的EventSource是最省事的方案但有两个限制一是只能接收 GET 请求二是无法自定义请求头。如果你的 SSE 接口需要鉴权 token且 token 放在 Authorization 头里那么EventSource直接歇菜。要么把 token 放 query 里有安全隐患但简单要么放弃EventSource用 fetch ReadableStream手动解析。我实际项目中使用的是 fetch 方案因为我们的鉴权更严格。核心代码如下const response await fetch(/api/agent/stream?question encodeURIComponent(question), { headers: { Authorization: Bearer ${token}, }, }); const reader response.body!.getReader(); const decoder new TextDecoder(utf-8); let buffer ; while (true) { const { value, done } await reader.read(); if (done) break; buffer decoder.decode(value, { stream: true }); const lines buffer.split(\n); buffer lines.pop()!; for (const line of lines) { if (line.startsWith(event:)) { currentEvent line.slice(6).trim(); } else if (line.startsWith(data:)) { const data JSON.parse(line.slice(5).trim()); handleEvent(currentEvent, data); } } }这段代码的要点是缓冲区边界处理TCP 流是不可控的一个完整事件可能被拆成两次read()返回所以必须维护一个buffer。split(\n)之后保留最后一段残片等下一次读取再拼上车。这个细节不做你在网络稍微波动时就会遇到“JSON 解析错误”的玄学 bug。如果你只是做内部工具token 放 query 也能接受那直接用EventSource就行连上面这些解析代码都省了。取舍取决于你的安全要求。5. 我没躲过去的那些坑超时、断连、鉴权5.1 before completion: idle timeout waiting for sse 的完整排查链路这个报错我在线上吃过一次亏。现象是某些长回答的请求会在大约 30 秒后中断前端 EventSource 自动重连但重连后 Agent 重新跑了一遍回答就重复了。排查过程值得完整记录下来。第一步看 Nginx 配置。我起初怀疑是proxy_read_timeout但这个值我设的是 120s报错却 30s 就出现对不上。把 Nginx 日志调到 debug 级别发现连接并不是 Nginx 掐断的而是上游 Node.js 进程主动断的。第二步看 NestJS 日志。每次报错前最后一条日志都是 LLM 调用开始之后就没有任何输出直到报Error: Idle timeout waiting for SSE。去查了nestjs/platform-express的源码发现它的 SSE 实现里有一个隐藏机制如果Observable超过某个时间没有发射数据框架会认为连接空闲主动结束响应。第三步确认根因。空闲超时的判断逻辑是“事件间隔”不是“总时长”。LLM 在做一次工具调用或长上下文推理时经常出现 20~30 秒的静默期期间既没有 token 也没有任何事件。对于 SSE 连接来说这确实是“空闲”于是框架就把连接关了。根因清楚了解法就很简单加入定时心跳。我用interval(15000)定期发射心跳事件保证事件间隔永远小于空闲超时阈值。改完之后这个报错再没出现过。如果你用的不是 NestJS也建议检查一下你的 Web 框架对 SSE 空闲连接有没有类似的默认回收逻辑。5.2 用户断开连接时下游 LLM 调用还在跑SSE 连接断开后后端闭包里的for await循环未必感知得到。我在一次压测时发现用户刷新页面后后端对 OpenAI 的调用还在持续白白烧了几千个 token。排查后发现RxJSsubscriber有个closed属性当客户端断开时subscriber.closed会变成true。但在我的run()函数的内部循环里并没有检查这个状态。修正后的代码for await (const chunk of generator) { if (subscriber.closed) { // 客户端断了终止 Agent 执行 await generator.return?.(); break; } subscriber.next({ ... }); }我自己之前见过一个解决方案是用AbortController把 abort 信号传给 LangChain 的调用链路。LangChain 底层其实原生支持signal配置很多 HTTP 调用和流式请求都会主动监听它。NestJS 的Req() req: Request对象里有req.on(close, ...)事件你可以在请求关闭时触发abortController.abort()。这样下游的 HTTP 客户端如 OpenAI SDK会立刻中止请求效率比轮询subscriber.closed高很多。我是两者都做了循环里检查closed闭包注册close事件触发abort()双保险。5.3 SSE 鉴权的另一个隐蔽问题守卫与请求头NestJS 的Sse()装饰器本质是一个普通 GET 接口所以你可以正常使用UseGuards(AuthGuard)来鉴权。守卫在 handler 执行之前跑所以鉴权失败时根本不会建立 SSE 连接直接返回 401。但是这里有一个隐蔽的问题某些前端库或者EventSource默认的Accept头是*/*而 NestJS 的 SSE 端点如果启用了全局的 Content-Type 校验比如Accept: application/json拦截器会导致连接被提前拒绝。解决方案是给 SSE 端点加白名单或者要求前端显式传入Accept: text/event-stream。还有一个更隐蔽的是token 放在 Authorization 头里时某些代理服务器会把 Authorization 头记录到访问日志中造成 token 泄露风险。如果对安全要求很高建议做一个“一次性 SSE 票据”方案先用普通 API 换一个短时效的 streamTicketSSE 请求时用它作为鉴权凭证。这样即便日志泄露票据很快失效影响面可控。6. 从 Demo 到稳定上线我做的三个优化6.1 跨实例推送用 Redis Pub/Sub 解决多节点问题单实例跑 SSE 没有任何问题但一旦服务横向扩展到多节点就出现一个经典困境Nginx 把 SSE 请求负载均衡到了节点 A但产生事件的 Agent 任务跑在了节点 B 上。解决办法是引入一个消息总线让所有节点订阅同一个频道收到事件后再推给对应的 SSE 连接。我用的是 Redis Pub/Sub。具体流程是NestJS 收到 SSE 请求后生成一个唯一的streamId把Observable的 subscriber 注册到一个本地 Map 里Agent 任务在任意节点执行时每产出一个事件就redis.publish(streamId, JSON.stringify(event))所有节点都订阅同一个 channel收到消息后通过streamId查找本地的 subscriber如果命中了就推出去。这个方案的重价在于Map 只存在于单个节点上发布和订阅必须配对。你可能要稍微改改架构比如建立连接时既订阅也发布或者用 Redis Stream 做持久化。如果只是中小规模也可以直接把所有 SSE 连接都接到同一个节点上用 sticky session省下大量复杂度。6.2 重连与补发用 Last-Event-ID 避免答案丢失SSE 断线重连后前端重新建立连接但服务端并不知道前端已经收到哪些事件。如果 Agent 已经执行到一半断了重连后从头跑一遍显然是浪费而且用户看到答案重复。我选了协议自带的Last-Event-ID机制后端每个事件附带递增 ID前端在断线重连时浏览器或者 fetch 代码会把上一次收到的 ID 带给后端。后端拿到Last-Event-ID后从 Redis 里查询该 ID 之后的事件并补发。这个方案有一个前提事件需要持久化。我是用 Redis List 存储最近一小时的 Agent 事件每次连接就把当前的streamId和已发送事件的 ID 范围记录下来。补发逻辑虽然不复杂但实现时注意清理过期数据不然 Redis 内存会涨得很快。6.3 不要忽略了 Observability监测 SSE 连接和事件时延上线后我发现一个问题用户体感忽快忽慢但后端平均响应时间看起来很正常。原因在于平均时间掩盖了长尾延迟。后来我在事件归一化层加了一些打点统计每次agent:message事件从产出到推给客户端的耗时以及工具调用的耗时分布。这边还有一个踩过的坑SSE 连接不主动关闭TCP 连接会一直占着。如果用户打开页面后长时间不管连接积累多了会打满文件描述符。所以我在前端加了一个页面visibilitychange监听页面隐藏超过 5 分钟就主动断开 SSE 连接后端也设置一个最大连接时长比如 30 分钟强制结算。另外一个数据处理细节SSE 的流式事件是高频的如果直接全量接日志系统会炸掉日志索引。我会对事件做采样比如agent:message只统计长度、时间和首字节时延只在出现错误时才记录完整内容。7. 写在最后建议从这张图开始落地你的 SSE Agent如果你现在要动手做类似的需求我建议从最简版本开始先写一个后端Sse()端点返回固定事件前端用fetch流式解析打出来确认 SSE 通路完全正常再把 LangChain 的流式迭代器接进去最后才是心跳、鉴权、断连重连这些生产级细节。不要一开始就上 Redis、多节点、补发机制大概率会把自己绕晕。最后送一个小技巧调试 SSE 时直接在命令行用 curl 看原始返回比打开浏览器 DevTools 更直观。curl -N http://localhost:3000/api/agent/stream?questionxxx-N参数是关键它关闭了 curl 的缓冲你就能看到事件一条条蹦出来了。配合jq解析 JSON 事件排查问题速度翻倍。这段链路我现在用得滚瓜烂熟希望也能帮你省下几个晚上的排查时间。