三亩地 三亩地SAN MU DI · CODE DIARY
ARTICLE DETAIL

日记详情

真实记录编程学习的某一天,欢迎挑你感兴趣的翻一翻。

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

基于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): Promise<string> { // 调用真实天气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的以下几个部分:

  1. ChatModel: 我们使用ChatOpenAIChatAnthropic等类来与LLM对话。关键是要配置streaming: true以启用流式输出。
  2. Tools: 将我们自己的业务功能(如WeatherService)包装成LangChain能识别的工具。这需要定义一个name,description(非常重要,Agent靠这个描述决定是否调用该工具)和schema(输入参数的定义)。
  3. AgentExecutor: 这是Agent的大脑。它接收一个AgentType(如OPENAI_FUNCTIONS)和定义好的工具列表,并提供一个invokestream方法来运行。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): Promise<string> { // 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做两件主要事情:

  1. 将LangChain Agent的stream方法返回的迭代器,转换成一个Observable流。
  2. 在这个流中插入自定义逻辑,比如解析工具调用事件、将工具执行结果格式化后重新喂给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): Observable<any> { 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.ts

3.1 AI模块(AiModule)的职责

AiModule是AI功能的入口模块。它的主要职责是:

  • 导入ToolsModuleAgentsModule
  • 通过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的Agent?LangChain的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)的实现

这是最核心的部分。这个服务将:

  1. 注入配置好的ChatModelToolRegistryService
  2. 创建LangChain的AgentExecutor
  3. 暴露一个公共方法(如streamResponse),接收用户查询,返回一个Observable流。

创建AgentExecutor

import { 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 Observable

AgentExecutor.stream()方法返回一个AsyncIterable。我们需要将其转换为一个能发出丰富事件的Observable

// streaming.agent.service.ts 中的核心方法 import { Observable, from } from 'rxjs'; import { map, catchError } from 'rxjs/operators'; streamResponse(userInput: string): Observable<StreamEvent> { // 定义一个自定义事件类型 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 Observable<StreamEvent>((subscriber) => { (async () => { try { const stream = await this.agentExecutor.stream({ input: userInput, }); for await (const chunk of stream) { // chunk 的结构取决于Agent类型和LangChain版本 // 对于OpenAI Functions Agent,chunk可能包含: // - `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控制器中暴露流式端点

现在,我们可以在控制器中创建一个SSE(Server-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): Observable<MessageEvent> { 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支持实时搜索建议,可以使用debounceTimedistinctUntilChanged来避免频繁请求。
  • 错误恢复:使用retrycatchError操作符,当工具调用失败时,尝试备用方案或给出友好提示。
  • 流的组合:如果回答需要综合多个数据源,可以使用forkJoincombineLatest来并行执行多个工具调用,然后合并结果流。

示例:在工具调用时添加加载状态和超时

import { timeout, catchError } from 'rxjs/operators'; // 在streamResponse方法内部的事件处理中,对工具调用事件进行处理 // 假设我们有一个执行工具的方法,它返回一个Observable executeTool(toolName: string, input: any): Observable<string> { 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_starttool_endtoken事件。上面的示例给出了一个基本框架,但实际实现需要你深入理解所选AgentExecutor的流输出结构,并可能需要对LangChain的事件回调(callbacks)进行更底层的定制。

5. 错误处理、监控与性能优化

一个生产级的Agent必须健壮。以下是一些关键考量点。

5.1 结构化错误处理

错误可能来自多个层面:

  1. LLM API错误:如网络超时、额度不足、模型过载。需要在LlmProvider或调用处设置重试逻辑和友好的降级提示。
  2. 工具执行错误:如数据库查询失败、第三方API不可用。工具自身应捕获异常并返回格式化的错误信息给Agent,而不是抛出异常导致整个流崩溃。Agent应该能理解工具返回的错误并做出反应(如“我无法获取天气信息,请稍后再试”)。
  3. Agent逻辑错误:如陷入循环、产生不符合预期的输出。可以设置最大迭代次数(maxIterations)来强制停止。

在RxJS流中,务必使用catchError操作符来捕获错误,并向下游发出一个格式化的错误事件,让前端能够展示,而不是让流无声无息地终止。

5.2 日志与监控

使用NestJS内置的Logger或集成像Winston这样的日志库,对关键事件进行记录:

  • 用户查询内容(注意隐私,可脱敏)。
  • 调用的工具及其参数。
  • 工具执行耗时。
  • 最终响应Token数量。
  • 发生的任何错误。

这对于调试、分析Agent行为、计算成本至关重要。可以考虑将日志结构化后输出到stdout,然后由日志收集系统(如ELK)处理。

5.3 性能与成本优化

  • 流式传输:本身就能极大提升用户体验感知性能。
  • 缓存:对于一些耗时的工具调用结果(如相对稳定的数据查询),可以考虑在服务层添加缓存(如Redis),避免重复调用。但要注意缓存数据的时效性。
  • Token管理:在系统提示词中明确要求Agent回答简洁。监控每次交互的输入+输出Token数,设置上限,防止恶意或意外的长文本消耗。
  • 连接管理:对于SSE/WebSocket连接,要做好心跳和超时管理,及时释放资源。

6. 测试策略:如何测试一个流式AI Agent?

测试是保证复杂系统稳定性的关键。我们需要分层测试:

  1. 单元测试(Unit Test)

    • 工具测试:单独测试每个工具类,模拟其依赖(如HTTP服务),验证给定输入能否产生正确输出。
    • 服务逻辑测试:测试StreamingAgentService中不涉及LangChain和RxJS流的纯逻辑部分。可以使用Jest等框架。
  2. 集成测试(Integration Test)

    • Agent流程测试:使用模拟的LLM(如ChatOpenAIcall方法可以被Jest Mock)和模拟的工具,测试整个AgentExecutor的调用流程。验证在给定输入下,是否会触发预期的工具调用序列。
    • HTTP端点测试:使用supertest测试/agent/stream端点,验证其是否能正确建立SSE连接并返回预期格式的事件流。这里可以模拟StreamingAgentService返回一个固定的Observable序列。
  3. E2E测试(End-to-End Test)

    • 在接近生产的环境(使用测试环境的API密钥)中,运行一组代表性的用户查询,验证从请求到最终流式输出的完整流程是否符合预期。这类测试运行较慢且可能有成本,适合在CI/CD的关键节点运行。

测试RxJS流的心得:测试Observable可以使用rxjsTestScheduler,但学习曲线较陡。一个更实用的方法是:在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中注册它,整个流式交互的框架就能自动适应,这才是“可扩展性”的真正体现。

← 返回列表