拓冰建站拓冰建站
首页 / 资讯中心 / 正文

基于NestJS、LangChain与RxJS构建可扩展AI智能体架构实践

1. 从单体AI应用到可扩展Agent架构的演进最近在折腾一个内部的知识库问答系统最初版本很简单一个FastAPI接口里面塞了个LangChain的链用户提问模型回答完事。但随着需求越来越复杂比如要支持多轮对话、实时流式输出、调用外部API查天气、查数据库甚至还要根据对话历史动态决定下一步动作原来的“面条式”代码很快就变成了“意大利面式”的灾难。每次加新功能都感觉在拆东墙补西墙测试更是噩梦。这让我意识到当AI应用从简单的“一问一答”升级为具备自主决策和行动能力的“智能体”时架构的选型至关重要。这也是为什么我开始深入研究并最终选择用NestJS、LangChain和RxJS这套组合拳来重构我的AI Agent。简单来说我们想打造的不是一个只会回话的聊天机器人而是一个可扩展、可观测、能流式响应且具备工具调用能力的AI智能体。这里的“可扩展”意味着新工具、新模型、新逻辑能像乐高积木一样轻松插拔“可观测”是指我们能清晰追踪Agent的每一次思考、每一次工具调用的输入输出“流式响应”则是为了提供类似ChatGPT那样逐字输出的实时体验降低用户等待的焦虑感而“工具调用”是Agent智能的核心让它能从“知道”升级到“能做到”。为什么是这三个技术栈NestJS提供了企业级Node.js后端所需的模块化、依赖注入和清晰的生命周期管理让复杂的业务逻辑变得井然有序。LangChain是AI应用开发的“瑞士军刀”其Agent和Tool抽象完美契合我们的需求。而RxJS这个强大的响应式编程库则是处理流式数据、管理异步事件和构建复杂数据管道的“神兵利器”。三者结合恰好能解决构建复杂AI Agent时遇到的架构混乱、状态管理困难和实时性挑战。2. 技术栈深度解析为何是NestJS、LangChain与RxJS在决定技术栈之前我对比过不少方案。比如用纯Express/ FastAPI搭个简单服务或者用更轻量的框架。但最终选择这三者是经过深思熟虑的它们各自解决了AI Agent工程化中的不同痛点。2.1 NestJS为复杂AI Agent提供坚实底座很多人觉得NestJS重但对于一个需要长期维护和迭代的AI Agent项目它的“重”恰恰是优势。AI Agent的内部状态、工具注册、模型配置、对话历史管理这些都不是简单的对象而是有明确生命周期和依赖关系的组件。NestJS的依赖注入容器能优雅地管理这些组件。例如我可以定义一个ToolRegistryService所有工具都在这里注册AgentService依赖这个注册表来获取可用工具而AgentController只负责接收HTTP请求并调用AgentService。这种清晰的边界让代码的可测试性和可维护性大大提升。更重要的是NestJS的模块化思想。我可以把LangChain相关的配置模型、提示词、记忆放在一个独立的LangChainModule里把工具的实现放在ToolsModule里把处理流式响应的逻辑放在StreamingModule里。当需要新增一个“查询股票价格”的工具时我只需要在ToolsModule中新增一个Service并在ToolRegistryService中注册即可其他模块完全不受影响。这种架构对于需要频繁迭代和扩展的AI应用来说是至关重要的。2.2 LangChainAI智能体的“大脑”与“工具箱”LangChain的核心价值在于它提供了一套高层次的抽象让我们不必从零开始处理与大模型交互的复杂性。对于Agent而言最关键的两个抽象是AgentExecutor和Tool。AgentExecutor是驱动Agent运行的核心引擎。它接收用户输入、可用工具列表和对话历史然后根据预设的Agent类型如OpenAI Functions, ReAct来协调整个过程调用模型进行思考、决定是否使用工具、解析工具调用参数、执行工具、将工具结果返回给模型进行下一步思考直到模型给出最终答案。这个循环逻辑如果自己实现会非常繁琐且容易出错。Tool抽象则统一了外部能力的接口。无论是调用一个HTTP API、执行一段数据库查询还是运行一段Python代码都可以包装成一个Tool。LangChain负责将工具的description和parameters以模型能理解的格式如JSON Schema提供给大模型并负责解析模型返回的工具调用指令。这让我们能够以声明式的方式为Agent“装配”各种能力。2.3 RxJS驾驭流式数据与复杂异步流这是整个架构中最为精妙也最具挑战的一环。传统的请求-响应模式在Agent场景下会卡住因为Agent的思考-行动循环可能需要多步每一步都可能耗时。我们需要的是像WebSocket或SSE那样能够持续向客户端推送中间状态和最终结果的能力。RxJS的Observable可观察对象是描述这种异步数据流的完美模型。我们可以将一次Agent的执行过程建模为一个数据流流开始用户输入事件。流中数据模型思考的中间文本“我需要调用天气工具...”、工具调用的开始和结束事件、工具执行的结果、模型根据结果生成的后续文本。流结束Agent返回最终答案。使用RxJS的Operators操作符我们可以像组装流水线一样处理这个流过滤掉调试信息、将工具调用结果格式化为更友好的展示、在流结束时自动清理资源、甚至处理多个并发Agent请求之间的背压问题。它让复杂的异步逻辑变得声明式和可组合这是用回调或Promise链难以优雅实现的。三者之间的关系可以这样比喻NestJS是身体和骨架提供了结构和组织LangChain是大脑和神经系统负责思考和发出指令RxJS是血液循环系统负责在各个环节之间高效、有序地传输信息和状态。缺少任何一个构建一个健壮、可扩展的流式Agent都会变得异常困难。3. 核心架构设计与模块拆解有了技术栈的认知我们来搭建这个Agent的骨架。我的目标是设计一个松耦合、高内聚的架构让数据流清晰可见。下面是我在实践中总结出的核心模块设计。3.1 项目结构与模块划分一个典型的NestJS项目结构如下我们围绕Agent核心能力进行组织src/ ├── agent/ │ ├── agent.controller.ts # HTTP入口处理流式请求 │ ├── agent.service.ts # Agent核心协调逻辑集成LangChain与RxJS │ ├── agent.module.ts │ └── dto/ # 数据传输对象如提问请求体 ├── tools/ # 工具模块 │ ├── tools.module.ts │ ├── tool.registry.service.ts # 全局工具注册中心 │ ├── weather.tool.service.ts # 具体工具实现查天气 │ ├── calculator.tool.service.ts # 具体工具实现计算器 │ └── ... # 其他工具 ├── llm/ # 大语言模型配置模块 │ ├── llm.module.ts │ ├── llm-config.service.ts # 模型API密钥、基础URL等配置 │ └── langchain.provider.ts # 提供配置好的LangChain LLM实例 ├── streaming/ # 流式输出处理模块 │ ├── streaming.module.ts │ └── sse.adapter.ts # 服务器发送事件适配器 └── common/ # 公共装饰器、过滤器、拦截器3.2 Agent Service核心协调器的实现AgentService是整个系统的心脏。它不直接处理HTTP也不直接实现工具而是负责组装LangChain的Agent并用RxJS管道驱动执行流程。// agent.service.ts import { Injectable } from nestjs/common; import { Observable, from, mergeMap, concatMap, scan } from rxjs; import { AgentExecutor } from langchain/agents; import { Tool } from langchain/tools; import { LLMStreamingHandler } from ./llm-streaming.handler; Injectable() export class AgentService { constructor( private readonly toolRegistry: ToolRegistryService, private readonly llmService: LlmConfigService, ) {} async createAgentExecutor(): PromiseAgentExecutor { const tools: Tool[] this.toolRegistry.getAllTools(); const llm this.llmService.getChatModel(); // 获取配置好的ChatOpenAI等实例 // 使用OpenAI Functions Agent const agent await createOpenAIFunctionsAgent({ llm, tools, prompt: YOUR_CUSTOM_PROMPT, // 可以自定义提示词设定Agent角色 }); return AgentExecutor.fromAgentAndTools({ agent, tools, returnIntermediateSteps: true, // 关键返回中间步骤用于流式输出 }); } streamAgentExecution(userInput: string, sessionId: string): ObservableStreamEvent { // 1. 创建Agent执行器 const executorPromise this.createAgentExecutor(); // 2. 将LangChain的执行过程包装为Observable return from(executorPromise).pipe( concatMap(executor { // 关键使用callbacks来捕获流式输出和中间步骤 const streamHandler new LLMStreamingHandler(); return from( executor.invoke( { input: userInput }, { callbacks: [streamHandler] } // 注入回调处理器 ) ).pipe( // 与回调处理器产生的事件流合并 mergeMap(() streamHandler.getEventStream()), // 累积事件构建完整的响应上下文可选用于前端显示思考过程 scan((acc: StreamEvent[], event: StreamEvent) [...acc, event], []) ); }) ); } } // 定义流式事件类型 export type StreamEvent | { type: llm_new_token; token: string; } | { type: tool_start; toolName: string; input: string; } | { type: tool_end; toolName: string; output: string; } | { type: agent_finish; finalOutput: string; };这里的精髓在于LLMStreamingHandler它是一个自定义的LangChainCallbackHandler。LangChain在执行过程中会在关键节点模型生成新token、工具开始/结束执行等调用回调函数。我们在这个Handler里不是简单地打印日志而是将这些事件推入一个RxJS的Subject一种特殊的Observable从而将LangChain的内部过程转换为我们可控的流式事件源。3.3 工具模块能力插拔的基石工具模块的设计目标是实现“热插拔”。ToolRegistryService作为一个中央注册表所有工具服务在初始化时都向它注册。// tool.registry.service.ts import { Injectable, OnModuleInit } from nestjs/common; import { Tool } from langchain/tools; import { WeatherToolService } from ./weather.tool.service; Injectable() export class ToolRegistryService implements OnModuleInit { private tools new Mapstring, Tool(); constructor(private readonly weatherTool: WeatherToolService) {} onModuleInit() { this.registerTool(this.weatherTool.getTool()); // 注册其他工具... } registerTool(tool: Tool) { this.tools.set(tool.name, tool); } getAllTools(): Tool[] { return Array.from(this.tools.values()); } getTool(name: string): Tool | undefined { return this.tools.get(name); } }而具体的工具如WeatherToolService则负责将业务逻辑包装成LangChain能识别的Tool对象。这里的一个关键点是错误处理工具执行必须健壮并将任何错误信息以模型能理解的方式返回以便Agent能进行后续决策。// weather.tool.service.ts import { Injectable } from nestjs/common; import { Tool } from langchain/tools; import { z } from zod; // 用于参数验证 import { callWeatherAPI } from ./weather.api; // 假设的天气API客户端 Injectable() export class WeatherToolService { getTool(): Tool { const schema z.object({ location: z.string().describe(The city and state, e.g. San Francisco, CA), }); return new DynamicTool({ name: get_current_weather, description: Get the current weather in a given location. Input should be a location string., func: async (input: string) { try { // 解析输入参数 const { location } schema.parse(JSON.parse(input)); // 调用真实API const weather await callWeatherAPI(location); return The current weather in ${location} is ${weather.temperature}°C, ${weather.condition}.; } catch (error) { // 非常重要返回清晰的错误信息帮助Agent调整策略 return Failed to get weather for ${input}. Error: ${error.message}. Please ensure the location is valid.; } }, }); } }4. 流式响应与SSE集成实战流式响应是提升用户体验的关键。我们使用服务器发送事件技术来实现。NestJS控制器负责建立SSE连接并将AgentService返回的RxJS流映射为SSE格式的数据流。4.1 控制器与SSE适配器// agent.controller.ts import { Controller, Post, Body, Sse, MessageEvent } from nestjs/common; import { Observable, map } from rxjs; import { AgentService, StreamEvent } from ./agent.service; Controller(agent) export class AgentController { constructor(private readonly agentService: AgentService) {} Post(stream) Sse() // NestJS的SSE装饰器 handleStreamingRequest( Body() dto: { question: string; sessionId?: string } ): ObservableMessageEvent { const { question, sessionId default-session } dto; // 获取Agent执行的事件流 const event$ this.agentService.streamAgentExecution(question, sessionId); // 将内部事件流转换为SSE标准格式 return event$.pipe( map((events: StreamEvent[]) { // 这里可以根据业务需求决定每次发送单个事件还是批量事件 const latestEvent events[events.length - 1]; const data this.formatEventForSSE(latestEvent); return { data }; }) ); } private formatEventForSSE(event: StreamEvent): string { // 将事件对象序列化为JSON字符串前端可以解析并渲染 return JSON.stringify(event); } }4.2 前端如何消费这个流前端可以使用EventSourceAPI来轻松连接这个SSE端点。// 前端示例代码 const eventSource new EventSource(/agent/stream?question北京天气怎么样); eventSource.onmessage (event) { const data JSON.parse(event.data); switch (data.type) { case llm_new_token: // 将token追加到答案区域实现打字机效果 answerEl.innerText data.token; break; case tool_start: // 在UI上显示“正在查询天气...” showToolCallStatus(调用工具: ${data.toolName}); break; case tool_end: // 显示工具调用结果 showToolResult(data.output); break; case agent_finish: // 最终答案已生成可以关闭连接或进行其他操作 eventSource.close(); break; } }; eventSource.onerror (err) { console.error(SSE连接错误:, err); eventSource.close(); };4.3 处理背压与连接管理在实际生产中需要考虑更多细节。比如当客户端网络较慢时服务器端持续推送数据可能会导致背压。RxJS提供了丰富的操作符来处理例如使用bufferTime将一段时间内的事件缓冲后一次性发送而不是每个token都发一次。另外需要管理SSE连接的生命周期在客户端断开或Agent执行完毕时及时清理RxJS订阅避免内存泄漏。这可以通过在控制器方法中返回的Observable上使用finalize操作符来实现。5. 高级话题错误处理、可观测性与性能优化当基础功能跑通后下一个阶段就是让这个Agent系统变得健壮、透明和高效。这部分是区分玩具项目和生产级系统的关键。5.1 全方位的错误处理策略AI应用的错误处理非常复杂因为错误可能发生在模型调用、工具执行、网络传输等多个环节。模型调用错误如OpenAI API超时、配额不足、内容过滤。我们必须在LLMConfigService中设置重试逻辑使用指数退避和友好的降级提示“服务暂时不可用请稍后再试”。工具执行错误如前文所述工具内部必须用try-catch包裹返回结构化错误信息。更进一步可以在ToolRegistryService层面增加一个全局的工具调用拦截器对所有工具的执行进行监控和统一的错误日志记录。流传输错误SSE连接可能意外中断。我们需要在服务端感知连接状态一旦断开就立即停止Agent执行并清理资源避免做无用功。这可以通过监听RxJS流的complete和error事件以及NestJS的OnModuleDestroy生命周期钩子来实现。用户输入错误Agent可能被用户诱导执行不合理或有害的操作。除了在模型层面使用系统提示词进行约束外还应在工具层面进行输入验证和权限控制。例如一个“发送邮件”的工具绝不能允许用户通过Agent向任意地址发送任意内容。5.2 构建可观测性追踪Agent的思考过程对于调试和优化Agent可观测性至关重要。我们需要知道Agent在每一步“想”了什么为什么选择某个工具工具执行结果如何。结构化日志不要只用console.log。使用像Winston或Pino这样的日志库将每一次Agent执行分配一个唯一的correlationId然后记录下所有关键事件agent_start,llm_call,tool_selection,tool_execution,agent_finish以及它们的输入输出。这些日志应该输出到像ELK或Loki这样的集中式日志系统中。LangChain Callbacks我们已经用回调来驱动流式响应了同样可以利用它来记录更详细的信息。可以创建多个CallbackHandler一个用于流式输出一个专门用于记录结构化日志到文件或数据库。追踪与监控集成OpenTelemetry这样的分布式追踪系统。为一次完整的Agent调用创建一个Trace将模型调用、每个工具调用都作为Span。这样你可以在Jaeger或Zipkin的UI上直观地看到一次查询的时间都花在了哪里是模型慢还是某个工具慢。5.3 性能优化与缓存策略随着用户量增长性能问题会浮现。对话历史管理LangChain的BufferMemory会保存所有历史对话如果对话很长每次都会作为上下文发送给模型导致token消耗巨大、速度变慢且成本激增。解决方案是使用SummaryBufferMemory或自定义记忆类定期对历史对话进行摘要只保留关键信息。对于超长文档问答则需要引入RAG的检索步骤而不是把整个文档塞进上下文。工具调用缓存很多工具调用是幂等的比如查询某个城市当前的天气在短时间内结果不会变化。可以为工具层添加缓存例如使用Redis缓存键可以是工具名和输入参数的哈希。设置一个合理的TTL例如10分钟可以大幅减少对下游API的调用提升响应速度并降低成本。模型响应缓存对于常见或重复的问题可以直接缓存模型的最终输出。这需要谨慎因为同样的用户问题在不同上下文中可能有不同答案。一个更精细的策略是缓存“思考过程”即Agent的中间步骤如果遇到相同的问题和相似的上下文可以跳过部分步骤。异步与并行化如果Agent需要调用多个不相关的工具可以考虑在安全的前提下并行执行。RxJS的forkJoin或merge操作符可以帮助管理并行流。但要注意工具之间的依赖关系有依赖的工具必须按顺序执行。构建一个可扩展的AI流式Agent是一个系统工程它远不止是调用一个API那么简单。它涉及到前后端协作、异步编程模型、复杂的状态管理和生产级的运维考量。NestJS、LangChain和RxJS这套组合分别从应用架构、AI抽象和数据处理三个维度提供了强大的支撑。从我的实践来看初期学习曲线确实较陡尤其是RxJS的响应式思维需要时间适应。但一旦掌握其带来的代码清晰度、可维护性和应对复杂场景的能力是传统方式难以比拟的。最重要的是这个架构是面向未来设计的当你有新的想法比如接入语音、增加一个复杂的多Agent协作流程时你会发现现有的模块和模式能够很自然地容纳这些扩展而不是推倒重来。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门