ARTICLE DETAIL

建站实战干货

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

基于NestJS与LangChain构建可扩展的AI流式Agent架构实践

2026/8/13 22:59:35 拓冰建站 浏览量
基于NestJS与LangChain构建可扩展的AI流式Agent架构实践 1. 项目概述为什么我们需要一个“可扩展的AI流式Agent”如果你正在构建一个需要与大型语言模型LLM深度集成的后端应用比如一个智能客服、一个代码助手或者一个复杂的业务流程自动化工具你很可能已经感受到了几个痛点。第一LLM的响应是“阻塞式”的用户得等它“思考”完一整段话才能看到结果体验很差。第二当AI需要调用外部工具比如查数据库、调用API、执行计算时代码很容易变得一团乱麻各种回调地狱和状态管理让人头疼。第三随着业务逻辑变复杂如何优雅地组织代码、处理错误、进行监控和测试成了一个巨大的挑战。这就是“用 NestJS LangChain RxJS 打造可扩展的 AI 流式 Agent”这个项目要解决的核心问题。它不是一个简单的“Hello World”式集成而是一个面向生产环境的、企业级的解决方案架构。简单来说我们想打造一个这样的AI智能体Agent它能理解用户意图在需要时自主调用我们预先定义好的工具Tool Calling并且整个过程的结果是“流式”返回给前端的就像看视频缓冲一样一个字一个字地出来响应极快。同时整个后端架构要足够健壮、可测试、易于维护和扩展。为什么是这三个技术栈的组合NestJS 提供了一个开箱即用的、模块化、依赖注入清晰的企业级Node.js框架它能让我们的应用结构非常清晰。LangChain 是目前最流行的LLM应用开发框架它抽象了与各种模型交互、构建链Chain和智能体Agent的复杂性尤其是其工具调用和智能体执行器Agent Executor的设计非常成熟。而 RxJS这个响应式编程库则是实现“流式”体验和复杂异步流程控制的秘密武器。它能把LLM的响应、工具调用的异步操作、甚至是多个并行的AI任务都转换成可观测的数据流Observable让我们可以用声明式的方式组合、转换和控制这些流完美解决异步混乱的问题。这个项目适合已经对Node.js和TypeScript有基本了解并且希望将AI能力深度、优雅地集成到后端服务中的开发者。接下来我会带你从零开始拆解每一个核心环节分享我在实际搭建过程中踩过的坑和总结的最佳实践。2. 技术栈深度解析与选型考量在动手写代码之前我们必须理解为什么是这三个技术以及它们在这个架构中分别扮演什么角色。选型不当后期重构的成本会非常高。2.1 NestJS不只是另一个Web框架很多人把NestJS看作一个加强了装饰器的Express这低估了它的价值。在这个AI Agent项目中NestJS的核心价值在于其依赖注入DI容器和模块化架构。为什么是依赖注入想象一下你的Agent需要调用一个“查询天气”的工具这个工具又依赖于一个配置了API密钥的HTTP服务。如果没有DI你可能需要在Agent的构造函数里手动new一个HTTP客户端然后传递密钥。这会导致代码紧耦合难以测试因为你无法轻松替换为Mock的HTTP客户端。在NestJS中你可以通过Injectable()装饰器将WeatherService声明为一个提供者Provider然后在需要使用它的AgentService的构造函数中直接注入。测试时你可以轻松地提供一个模拟的WeatherService。// 一个工具服务 Injectable() export class WeatherService { constructor(private readonly httpService: HttpService) {} async getWeather(city: string): Promisestring { // 调用真实天气API const response await this.httpService.get(https://api.weather.com/${city}).toPromise(); return The weather in ${city} is ${response.data.condition}.; } } // Agent服务中注入并使用 Injectable() export class AgentService { constructor(private readonly weatherService: WeatherService) {} // 依赖注入 async process(query: string) { // ... LangChain Agent逻辑中可以调用 this.weatherService } }模块化则允许我们将功能拆分。例如我们可以有一个AiModule专门管理所有LangChain相关的配置模型、提示词模板一个ToolsModule注册所有可用的工具一个AgentsModule定义不同的智能体。这种清晰的分层使得代码库在增长时依然可维护。实操心得不要把所有LangChain的代码都堆在一个Service里。利用NestJS的模块将模型配置、工具定义、Agent执行器分别放在不同的提供者中。这样当你想切换模型比如从OpenAI换成Azure OpenAI时只需要修改AiModule中的一个配置文件而不是在几十个文件中搜索API密钥。2.2 LangChain从Chain到Agent的进化LangChain的核心概念是“链”Chain即把对LLM的调用、数据处理、工具调用等环节链接起来。而“智能体”Agent是一种特殊的链它引入了“推理”能力根据用户输入和上下文自主决定下一步是直接回答还是调用某个工具。在这个项目中我们重点关注LangChain的以下几个部分ChatModel 我们使用ChatOpenAI或ChatAnthropic等类来与LLM对话。关键是要配置streaming: true以启用流式输出。Tools 将我们自己的业务功能如WeatherService包装成LangChain能识别的工具。这需要定义一个name,description非常重要Agent靠这个描述决定是否调用该工具和schema输入参数的定义。AgentExecutor 这是Agent的大脑。它接收一个AgentType如OPENAI_FUNCTIONS和定义好的工具列表并提供一个invoke或stream方法来运行。stream方法正是我们实现流式响应的基础。工具定义示例import { Tool } from langchain/core/tools; import { Injectable } from nestjs/common; Injectable() export class WeatherTool extends Tool { name get_current_weather; description Get the current weather in a given location. Input should be a city name.; constructor(private readonly weatherService: WeatherService) { super(); } protected async _call(arg: string): Promisestring { // arg 是LLM解析出来的参数比如 Beijing try { const result await this.weatherService.getWeather(arg); return result; } catch (error) { return Failed to get weather: ${error.message}; } } }注意事项工具的描述description是Agent能否正确调用的关键。描述必须清晰、准确说明工具的功能和输入格式。比如“Input should be a city name.”就明确告诉LLM应该传入一个城市名字符串。模糊的描述会导致LLM无法理解或错误调用。2.3 RxJS响应式编程掌控异步流这是实现“可扩展”和“流式”的关键。LLM的流式响应本质是一个异步数据流每个Token都是一个事件。工具调用可能涉及多个异步HTTP请求。RxJS的Observable可以完美地表示这些流。核心优势组合性 你可以将LLM的Token流、工具调用的日志流、甚至用户中断请求的信号流通过操作符如mergeMap,switchMap,catchError进行灵活组合。背压处理 当生产数据LLM生成Token的速度快于消费速度网络传输时RxJS有策略来处理避免内存溢出。取消订阅 如果用户在前端关闭了连接你可以通过取消订阅Observable来立即中断LLM的生成过程节省资源和费用。在这个项目中我们会用RxJS做两件主要事情将LangChain Agent的stream方法返回的迭代器转换成一个Observable流。在这个流中插入自定义逻辑比如解析工具调用事件、将工具执行结果格式化后重新喂给Agent并保持流的连续性。一个简单的流转换示例import { Observable, from } from rxjs; async function* agentStreamGenerator(query: string) { const stream await agentExecutor.stream({ input: query }); for await (const chunk of stream) { yield chunk; // chunk可能是AgentAction、AgentFinish或中间Token } } // 在NestJS Service中 getAgentStream(query: string): Observableany { return from(this.agentStreamGenerator(query)); // 将异步生成器转为Observable }踩坑记录直接使用from转换复杂的LangChain流可能会遇到问题因为流中的事件结构复杂。更好的做法是使用new Observable(subscriber)手动创建在订阅函数中控制异步迭代器的遍历并针对不同事件类型onLLMNewToken,onToolStart调用subscriber.next()这样我们对流的控制力更强也更容易添加自定义逻辑和错误处理。3. 项目架构设计与核心模块拆解有了技术栈的理解我们来设计一个清晰、可扩展的项目结构。一个典型的NestJS项目结构如下我们将AI Agent的能力融入其中src/ ├── ai/ │ ├── ai.module.ts # AI核心模块导入模型和工具 │ ├── models/ # 模型配置相关 │ │ └── llm.provider.ts # 提供配置好的ChatModel实例 │ ├── tools/ # 所有工具定义 │ │ ├── tools.module.ts │ │ ├── weather.tool.ts │ │ └── calculator.tool.ts │ └── agents/ # 智能体定义 │ ├── agents.module.ts │ ├── base.agent.ts # 抽象基类封装通用流处理逻辑 │ └── streaming.agent.service.ts # 具体的流式Agent服务 ├── app.module.ts └── main.ts3.1 AI模块AiModule的职责AiModule是AI功能的入口模块。它的主要职责是导入ToolsModule和AgentsModule。通过LlmProvider一个自定义Provider来创建和配置LangChain的ChatModel实例。这里会集中管理API密钥、基础URL、模型温度temperature、最大Token数等参数。使用Provider而不是直接在Service里写死配置便于环境隔离开发/生产和动态配置。llm.provider.ts示例import { Provider } from nestjs/common; import { ChatOpenAI } from langchain/openai; export const LlmProvider: Provider { provide: LANGCHAIN_CHAT_MODEL, // 使用字符串或自定义Token作为标识 useFactory: () { return new ChatOpenAI({ modelName: gpt-4, temperature: 0.2, streaming: true, // 关键启用流式 openAIApiKey: process.env.OPENAI_API_KEY, // 其他配置... }); }, };3.2 工具模块ToolsModule的设计每个工具都是一个独立的类继承自LangChain的Tool同时也是一个NestJS的Injectable()服务。这样它既可以被LangChain的Agent使用也可以在里面注入其他NestJS服务如HttpService,Repository来执行业务逻辑。tools.module.ts需要将所有工具类放在providers数组中并导出ToolsModule以便AiModule导入。关键点如何将NestJS管理的工具实例传递给LangChain的AgentLangChain的Agent在初始化时需要工具实例的数组。我们可以在一个Service比如ToolRegistryService中通过依赖注入收集所有工具然后提供一个getTools()方法。// tool-registry.service.ts Injectable() export class ToolRegistryService { private tools: Tool[] []; // 通过构造函数注入所有工具NestJS会自动处理 constructor( private readonly weatherTool: WeatherTool, private readonly calculatorTool: CalculatorTool, ) { this.tools [this.weatherTool, this.calculatorTool]; } getTools(): Tool[] { return this.tools; } }3.3 智能体服务StreamingAgentService的实现这是最核心的部分。这个服务将注入配置好的ChatModel和ToolRegistryService。创建LangChain的AgentExecutor。暴露一个公共方法如streamResponse接收用户查询返回一个Observable流。创建AgentExecutorimport { Injectable, Inject } from nestjs/common; import { AgentExecutor, createOpenAIFunctionsAgent } from langchain/agents; import { ChatOpenAI } from langchain/openai; Injectable() export class StreamingAgentService { private agentExecutor: AgentExecutor; constructor( Inject(LANGCHAIN_CHAT_MODEL) private readonly chatModel: ChatOpenAI, private readonly toolRegistry: ToolRegistryService, ) { this.initializeAgent(); } private async initializeAgent() { const tools this.toolRegistry.getTools(); // 1. 定义提示词。系统提示词至关重要它决定了Agent的角色和行为准则。 const systemPrompt You are a helpful assistant. ... Use tools when necessary.; // 2. 创建Agent const agent await createOpenAIFunctionsAgent({ llm: this.chatModel, tools, prompt: systemPrompt, }); // 3. 创建执行器 this.agentExecutor new AgentExecutor({ agent, tools, returnIntermediateSteps: true, // 重要返回中间步骤便于流式展示 }); } }4. 流式响应与RxJS深度集成实战现在来到最具挑战也最有价值的部分如何让AgentExecutor的流不仅输出最终的文本还能实时反映其“思考过程”比如“我正在调用XX工具”4.1 包装LangChain的流为RxJS ObservableAgentExecutor.stream()方法返回一个AsyncIterable。我们需要将其转换为一个能发出丰富事件的Observable。// streaming.agent.service.ts 中的核心方法 import { Observable, from } from rxjs; import { map, catchError } from rxjs/operators; streamResponse(userInput: string): ObservableStreamEvent { // 定义一个自定义事件类型 type StreamEvent | { type: token; token: string } | { type: tool_start; toolName: string; input: string } | { type: tool_end; result: string } | { type: error; error: string } | { type: end }; return new ObservableStreamEvent((subscriber) { (async () { try { const stream await this.agentExecutor.stream({ input: userInput, }); for await (const chunk of stream) { // chunk 的结构取决于Agent类型和LangChain版本 // 对于OpenAI Functions Agentchunk可能包含 // - output: 最终的输出Token // - intermediateSteps: 中间步骤工具调用 if (chunk.output) { // 发出文本Token subscriber.next({ type: token, token: chunk.output }); } // 处理工具调用事件需要根据实际chunk结构调整 if (chunk.intermediateSteps chunk.intermediateSteps.length 0) { const lastStep chunk.intermediateSteps[chunk.intermediateSteps.length - 1]; if (lastStep.action) { // 工具开始 subscriber.next({ type: tool_start, toolName: lastStep.action.tool, input: JSON.stringify(lastStep.action.toolInput) }); } if (lastStep.observation) { // 工具结束返回结果 subscriber.next({ type: tool_end, result: lastStep.observation }); } } } subscriber.next({ type: end }); subscriber.complete(); } catch (error) { subscriber.next({ type: error, error: error.message }); subscriber.complete(); } })(); }); }4.2 在NestJS控制器中暴露流式端点现在我们可以在控制器中创建一个SSEServer-Sent Events或WebSocket端点。这里以更通用的SSE为例因为它基于HTTP更简单。// agent.controller.ts import { Controller, Get, Query, Res, Sse } from nestjs/common; import { Response } from express; import { Observable, interval } from rxjs; import { map } from rxjs/operators; import { StreamingAgentService } from ./streaming.agent.service; Controller(agent) export class AgentController { constructor(private readonly agentService: StreamingAgentService) {} Get(stream) Sse() // NestJS的SSE装饰器 streamAgentResponse(Query(q) query: string): ObservableMessageEvent { return this.agentService.streamResponse(query).pipe( map((event) { // 将自定义事件转换为SSE要求的MessageEvent格式 return { data: event, } as MessageEvent; }), catchError((error) { return of({ data: { type: error, error: error.message }, } as MessageEvent); }) ); } }前端可以通过EventSourceAPI连接到/agent/stream?q你的问题并监听onmessage事件来实时接收Token和工具调用状态。4.3 使用RxJS操作符增强流处理RxJS的强大之处在于操作符。例如我们可以防抖与搜索如果Agent支持实时搜索建议可以使用debounceTime和distinctUntilChanged来避免频繁请求。错误恢复使用retry或catchError操作符当工具调用失败时尝试备用方案或给出友好提示。流的组合如果回答需要综合多个数据源可以使用forkJoin或combineLatest来并行执行多个工具调用然后合并结果流。示例在工具调用时添加加载状态和超时import { timeout, catchError } from rxjs/operators; // 在streamResponse方法内部的事件处理中对工具调用事件进行处理 // 假设我们有一个执行工具的方法它返回一个Observable executeTool(toolName: string, input: any): Observablestring { return from(this.findToolAndExecute(toolName, input)).pipe( timeout(10000), // 10秒超时 catchError(err of(Tool ${toolName} execution failed: ${err.message})) ); } // 然后在主Observable中使用switchMap切换到工具执行流 // 这是一个概念性代码实际集成需要更精细的控制核心技巧处理工具调用的流式反馈是个难点。理想情况是前端看到“正在调用天气API...”然后看到“北京天气晴朗25度”最后LLM基于这个结果继续生成“所以建议你穿短袖”。这需要我们的Observable流能交错发出tool_start、tool_end和token事件。上面的示例给出了一个基本框架但实际实现需要你深入理解所选AgentExecutor的流输出结构并可能需要对LangChain的事件回调callbacks进行更底层的定制。5. 错误处理、监控与性能优化一个生产级的Agent必须健壮。以下是一些关键考量点。5.1 结构化错误处理错误可能来自多个层面LLM API错误如网络超时、额度不足、模型过载。需要在LlmProvider或调用处设置重试逻辑和友好的降级提示。工具执行错误如数据库查询失败、第三方API不可用。工具自身应捕获异常并返回格式化的错误信息给Agent而不是抛出异常导致整个流崩溃。Agent应该能理解工具返回的错误并做出反应如“我无法获取天气信息请稍后再试”。Agent逻辑错误如陷入循环、产生不符合预期的输出。可以设置最大迭代次数maxIterations来强制停止。在RxJS流中务必使用catchError操作符来捕获错误并向下游发出一个格式化的错误事件让前端能够展示而不是让流无声无息地终止。5.2 日志与监控使用NestJS内置的Logger或集成像Winston这样的日志库对关键事件进行记录用户查询内容注意隐私可脱敏。调用的工具及其参数。工具执行耗时。最终响应Token数量。发生的任何错误。这对于调试、分析Agent行为、计算成本至关重要。可以考虑将日志结构化后输出到stdout然后由日志收集系统如ELK处理。5.3 性能与成本优化流式传输本身就能极大提升用户体验感知性能。缓存对于一些耗时的工具调用结果如相对稳定的数据查询可以考虑在服务层添加缓存如Redis避免重复调用。但要注意缓存数据的时效性。Token管理在系统提示词中明确要求Agent回答简洁。监控每次交互的输入输出Token数设置上限防止恶意或意外的长文本消耗。连接管理对于SSE/WebSocket连接要做好心跳和超时管理及时释放资源。6. 测试策略如何测试一个流式AI Agent测试是保证复杂系统稳定性的关键。我们需要分层测试单元测试Unit Test工具测试单独测试每个工具类模拟其依赖如HTTP服务验证给定输入能否产生正确输出。服务逻辑测试测试StreamingAgentService中不涉及LangChain和RxJS流的纯逻辑部分。可以使用Jest等框架。集成测试Integration TestAgent流程测试使用模拟的LLM如ChatOpenAI的call方法可以被Jest Mock和模拟的工具测试整个AgentExecutor的调用流程。验证在给定输入下是否会触发预期的工具调用序列。HTTP端点测试使用supertest测试/agent/stream端点验证其是否能正确建立SSE连接并返回预期格式的事件流。这里可以模拟StreamingAgentService返回一个固定的Observable序列。E2E测试End-to-End Test在接近生产的环境使用测试环境的API密钥中运行一组代表性的用户查询验证从请求到最终流式输出的完整流程是否符合预期。这类测试运行较慢且可能有成本适合在CI/CD的关键节点运行。测试RxJS流的心得测试Observable可以使用rxjs的TestScheduler但学习曲线较陡。一个更实用的方法是在Service的方法中返回Observable在测试中订阅它并将发出的值收集到一个数组中然后断言这个数组是否符合预期序列。// 示例测试一个简单的流服务 it(should emit a sequence of events, (done) { const expectedEvents [ { type: token, token: Hello }, { type: tool_start, toolName: search }, { type: end } ]; const receivedEvents: any[] []; service.getSimpleStream().subscribe({ next: (event) receivedEvents.push(event), complete: () { expect(receivedEvents).toEqual(expectedEvents); done(); } }); });7. 部署与扩展思考当你的Agent开发完成后部署到生产环境需要考虑环境变量所有API密钥、模型端点URL等敏感信息必须通过环境变量注入。进程管理使用pm2或容器编排如Kubernetes来管理Node.js进程确保其崩溃后能自动重启。水平扩展NestJS应用本身是无状态的可以轻松水平扩展。但需要注意如果使用了内存缓存或Session需要转移到外部存储如Redis。Agent的版本化当你更新了提示词或工具集最好通过API版本如/v1/agent/stream或配置开关来逐步灰度发布方便回滚和A/B测试。扩展方向多Agent系统可以定义多个具有不同专长如客服、编程、分析的Agent并由一个路由Agent根据用户问题类型进行分发。这可以利用NestJS的模块化轻松实现。记忆Memory为Agent添加对话记忆使其能记住上下文。LangChain提供了多种记忆方案可以集成到我们的流式架构中通常是将记忆状态作为每次调用的一部分传入agentExecutor.stream()。与LangGraph集成对于更复杂、有状态、多分支的工作流可以探索LangChain的LangGraph。它本质上是一个基于图的编排框架可以用更直观的方式定义Agent之间的协作流程其执行过程同样可以流式化。构建这样一个可扩展的AI流式Agent是一次充满挑战但也收获巨大的工程实践。它迫使你深入思考异步编程、软件架构、用户体验和AI能力的结合。从我的经验来看最大的价值不在于快速实现一个原型而在于建立了一个清晰、健壮、易于迭代的基础设施。当产品经理提出“能不能让AI在回答前先查一下用户的历史订单”这样的需求时你只需要在ToolsModule中新增一个OrderHistoryTool然后在合适的Agent中注册它整个流式交互的框架就能自动适应这才是“可扩展性”的真正体现。