多智能体系统构建与部署实战:从原型到生产级落地
时间:2026-08-15 | 作者:318050 | 阅读:0Agent工程化实战:从原型到生产级多智能体系统的构建与部署
一、为什么Agent工程化如此重要?
2025年,大语言模型(LLM)的推理能力已突破实用门槛,基于LLM的智能体(Agent)从Demo走向真实业务场景。然而,原型与生产之间隔着一条巨大的工程化鸿沟:
可靠性:LLM输出不确定,如何保证业务流程的稳定性?可观测性:Agent的每一次工具调用、每一步推理都可能是黑盒,如何追踪和调试?扩展性:单Agent能力有限,如何设计多智能体协作并动态伸缩?性能:LLM响应延迟高,如何异步化、缓存、流式处理?运维:模型版本更新、API限流、成本控制,如何自动化管理?本文将围绕上述痛点,从零构建一个生产级多智能体系统,涵盖设计、编码、测试、容器化、部署与监控全流程。
二、系统架构设计
我们设计一个智能客服 代码辅助的双Agent协作系统,包含以下角色:
Agent | 职责 | 工具 |
|---|---|---|
Planner | 分析用户问题,拆解任务,调度执行 | 任务分解、资源检索 |
Executor | 执行具体动作,如查询数据库、调用API、生成代码 | SQL执行器、代码解释器、向量检索 |
Reviewer | 校验执行结果,必要时回退或重试 | 结果校验、反思 |
整体架构如下:
代码语言:ja vascript复制用户请求 → FastAPI网关 → 任务队列(Celery) → Agent Orchestrator ├─ Planner Agent (LLM) ├─ Executor Agent (LLM Tools) └─ Reviewer Agent (LLM)↓状态存储(Redis) 日志(ELK) 监控(Prometheus)
关键设计决策:
异步解耦:通过Celery将长耗时Agent推理任务与HTTP请求分离,避免超时和资源阻塞。状态管理:使用Redis维护会话上下文和Agent执行状态,支持断点续传。容错机制:每个Agent步骤增加超时、重试和降级逻辑(如LLM不可用时返回缓存或默认答案)。可观测性:集成OpenTelemetry,全链路追踪Agent的推理链和工具调用。三、核心实现:基于LangChain的自定义Agent
我们使用LangChain 0.3 (已适配v1.0 API)构建自定义Agent,并接入OpenAI兼容模型(支持GPT-4o、Claude、Qwen等)。
3.1 自定义工具
代码语言:ja vascript复制# tools.pyfrom langchain.tools import BaseToolfrom pydantic import BaseModel, Fieldimport subprocessimport sqlite3import jsonclass SQLQueryInput(BaseModel):query: str = Field(description="SQL查询语句")class SQLExecutorTool(BaseTool):name: str = "sql_executor"description: str = "执行SQL查询并返回结果,仅限SELECT语句"args_schema: type[BaseModel] = SQLQueryInputdef _run(self, query: str) -> str:# 实际生产环境应使用连接池和只读账户conn = sqlite3.connect("data.db")try:cursor = conn.cursor()cursor.execute(query)rows = cursor.fetchall()columns = [desc[0] for desc in cursor.description]return json.dumps([dict(zip(columns, row)) for row in rows], ensure_ascii=False)except Exception as e:return f"SQL执行错误: {str(e)}"finally:conn.close()class CodeInterpreterTool(BaseTool):name: str = "python_interpreter"description: str = "执行Python代码并返回标准输出,用于数学计算或数据处理"def _run(self, code: str) -> str:try:# 使用subprocess隔离执行,限制资源result = subprocess.run(["python3", "-c", code],capture_output=True,text=True,timeout=10,env={"PYTHONPATH": ""})return result.stdout or result.stderrexcept subprocess.TimeoutExpired:return "代码执行超时(10秒)"except Exception as e:return f"执行异常: {str(e)}"
3.2 Agent定义与提示工程
每个Agent拥有独立的系统提示和工具集,我们采用ReAct风格提示,并加入结构化输出要求。
代码语言:ja vascript复制# agents.pyfrom langchain_openai import ChatOpenAIfrom langchain.agents import create_react_agent, AgentExecutorfrom langchain.prompts import PromptTemplatefrom tools import SQLExecutorTool, CodeInterpreterTool# 基础LLM(支持环境变量切换模型)llm = ChatOpenAI(model=os.getenv("LLM_MODEL", "gpt-4o-mini"),temperature=0.1,timeout=60,max_retries=3,)planner_prompt = PromptTemplate.from_template("""你是一个任务规划专家。根据用户输入,拆解为可执行的子任务,并决定由哪个Agent执行。用户输入: {input}已有上下文: {context}可用工具: {tools}请按JSON格式输出计划:{{"plan": [{{"step": 1, "agent": "executor", "action": "sql_query", "params": {{"query": "..."}}}},{{"step": 2, "agent": "reviewer", "action": "validate", "params": {{"expected": "..."}}}}]}}""")executor_prompt = PromptTemplate.from_template("""你是一个执行专家,负责调用工具完成具体任务。任务描述: {task}可用工具: {tools}请逐步思考并调用工具,最后给出执行结果。""")def create_planner_agent():tools = []# Planner本身不直接调用工具,而是输出计划# 但为了统一,我们使用工具调用方式,这里特殊处理# 实际使用中,可让Planner直接调用"任务分解"工具,但为了演示,我们直接解析JSONpass# 更实用的方式:使用langchain的structured outputfrom langchain.output_parsers import PydanticOutputParserfrom pydantic import BaseModel, Fieldfrom typing import Listclass PlanStep(BaseModel):step: intagent: straction: strparams: dictclass Plan(BaseModel):plan: List[PlanStep]parser = PydanticOutputParser(pydantic_object=Plan)planner_llm = llm.with_structured_output(Plan)# 注意:with_structured_output需要模型支持JSON模式,否则使用function calling# 实际执行Agentdef create_executor_agent():tools = [SQLExecutorTool(), CodeInterpreterTool()]agent = create_react_agent(llm, tools, executor_prompt)return AgentExecutor(agent=agent, tools=tools, verbose=True, handle_parsing_errors=True)executor = create_executor_agent()
3.3 多智能体编排器(Orchestrator)
编排器负责管理对话状态、调用各Agent、处理结果汇总。
代码语言:ja vascript复制# orchestrator.pyimport jsonimport asynciofrom typing import Dict, Anyfrom agents import planner_llm, executor, reviewer_llm# reviewer类似定义,省略from redis import Redisimport celeryredis_client = Redis(host='redis', decode_responses=True)class AgentOrchestrator:def __init__(self, session_id: str):self.session_id = session_idself.context = self._load_context()def _load_context(self) -> Dict:key = f"session:{self.session_id}:context"data = redis_client.get(key)return json.loads(data) if data else {"history": [], "last_plan": None}def _sa ve_context(self):key = f"session:{self.session_id}:context"redis_client.setex(key, 3600, json.dumps(self.context))async def process_request(self, user_input: str) -> str:# 1. 规划plan_prompt = f"用户输入: {user_input}上下文: {self.context['history'][-5:]}"try:plan = await planner_llm.ainvoke(plan_prompt)except Exception as e:# 降级:直接执行默认动作return self._fallback(user_input)self.context['last_plan'] = plan.dict()self._sa ve_context()# 2. 按计划执行(可并行执行独立步骤)results = []for step in plan.plan:if step.agent == "executor":result = await self._run_executor(step.action, step.params)elif step.agent == "reviewer":result = await self._run_reviewer(step.action, step.params, previous=results)# 其他agent...results.append({"step": step.step, "result": result})# 3. 汇总生成最终回复final = await self._generate_final_response(user_input, results)self.context['history'].append({"user": user_input, "assistant": final})self._sa ve_context()return finalasync def _run_executor(self, action: str, params: dict) -> str:# 映射action到具体工具调用if action == "sql_query":tool = SQLExecutorTool()return tool._run(params.get("query", ""))elif action == "python_code":tool = CodeInterpreterTool()return tool._run(params.get("code", ""))else:# 使用AgentExecutor通用处理response = await executor.ainvoke({"input": f"执行{action},参数{params}"})return response["output"]# reviewer和最终回复方法类似,使用LLM调用
四、工程化关键:异步任务与API层
为避免HTTP请求阻塞,我们将Agent执行放入Celery异步任务,并支持WebSocket流式返回。
4.1 Celery任务定义
代码语言:ja vascript复制# tasks.pyfrom celery import Celeryfrom orchestrator import AgentOrchestratorapp = Celery('agent_tasks', broker='redis://redis:6379/0')app.conf.update(task_serializer='json',accept_content=['json'],result_serializer='json',timezone='Asia/Shanghai',task_track_started=True,task_time_limit=300,# 5分钟超时task_soft_time_limit=240,task_acks_late=True,# 防止任务丢失)@app.task(bind=True, max_retries=3)def process_agent_request(self, session_id: str, user_input: str):try:orch = AgentOrchestrator(session_id)# 异步执行loop = asyncio.new_event_loop()asyncio.set_event_loop(loop)result = loop.run_until_complete(orch.process_request(user_input))loop.close()return resultexcept Exception as e:# 重试机制self.retry(exc=e, countdown=2 ** self.request.retries)
4.2 FastAPI接口(支持同步轮询和WebSocket流式)
代码语言:ja vascript复制# api.pyfrom fastapi import FastAPI, WebSocket, BackgroundTasksfrom pydantic import BaseModelfrom tasks import process_agent_requestfrom celery.result import AsyncResultimport jsonapp = FastAPI(title="Multi-Agent System")class RequestPayload(BaseModel):session_id: strmessage: str# 同步轮询接口(适用于短任务)@app.post("/agent/sync")async def sync_agent(payload: RequestPayload):task = process_agent_request.delay(payload.session_id, payload.message)# 阻塞等待(生产环境不建议,这里仅演示)result = task.get(timeout=60)return {"status": "success", "data": result}# 异步提交,轮询结果@app.post("/agent/async")async def async_agent(payload: RequestPayload):task = process_agent_request.delay(payload.session_id, payload.message)return {"task_id": task.id, "status_url": f"/agent/status/{task.id}"}@app.get("/agent/status/{task_id}")async def get_status(task_id: str):task = AsyncResult(task_id)if task.ready():return {"status": "completed", "result": task.result}elif task.failed():return {"status": "failed", "error": str(task.info)}else:return {"status": "pending", "progress": task.info.get('progress', 0)}# WebSocket流式返回(支持推理过程实时展示)@app.websocket("/ws/agent")async def websocket_agent(websocket: WebSocket):await websocket.accept()try:data = await websocket.receive_json()session_id = data.get("session_id")message = data.get("message")# 使用Celery的异步任务,但通过Redis pub/sub推送中间结果# 这里简化:直接调用orchestrator并流式生成from orchestrator import AgentOrchestratororch = AgentOrchestrator(session_id)# 模拟流式输出(实际可用async生成器)async for chunk in orch.stream_process(message):await websocket.send_text(json.dumps({"type": "chunk", "data": chunk}))await websocket.send_text(json.dumps({"type": "end"}))except Exception as e:await websocket.send_text(json.dumps({"type": "error", "data": str(e)}))finally:await websocket.close()
五、可观测性:全链路追踪与日志
使用OpenTelemetry Jaeger实现分布式追踪,每条请求生成唯一trace_id,贯穿网关、Celery、LLM调用、工具执行。
# tracing.pyfrom opentelemetry import tracefrom opentelemetry.exporter.jaeger.thrift import JaegerExporterfrom opentelemetry.instrumentation.fastapi import FastAPIInstrumentorfrom opentelemetry.instrumentation.celery import CeleryInstrumentorfrom opentelemetry.sdk.trace import TracerProviderfrom opentelemetry.sdk.trace.export import BatchSpanProcessordef setup_tracing(app):provider = TracerProvider()processor = BatchSpanProcessor(JaegerExporter(agent_host_name=os.getenv("JAEGER_HOST", "jaeger"),agent_port=6831,))provider.add_span_processor(processor)trace.set_tracer_provider(provider)FastAPIInstrumentor.instrument_app(app)CeleryInstrumentor().instrument()
日志结构化:使用structlog,每个日志条目包含trace_id、session_id、agent_name等。
来源:整理自互联网
免责声明:文中图文均来自网络,如有侵权请联系删除,心愿游戏发布此文仅为传递信息,不代表心愿游戏认同其观点或证实其描述。
相关文章
更多-
- 工作流驱动智能体架构的场景定义编排与执行解析
- 时间:2026-08-15
-
- 苹果Xcode 27集成AI智能体开启智能编程新时代
- 时间:2026-08-15
-
- 国内AI安全产品市场深度分析与产业升级趋势
- 时间:2026-08-15
-
- OpenWorker:吴恩达开源AI桌面智能体工具介绍
- 时间:2026-08-15
-
- 发布纳米Work企业智能体平台首批用户获赠1亿Token
- 时间:2026-08-14
-
- 发布纳米Work企业智能体平台 首批用户获1亿Token试用额度
- 时间:2026-08-13
-
- 心言集团推出跨智能体工作流终端AgentLink
- 时间:2026-08-13
-
- Arm在WAIC 2026揭示AI智能体落地两大趋势
- 时间:2026-08-13
精选合集
更多大家都在玩
热门话题
大家都在看
更多-
- 多智能体系统构建与部署实战:从原型到生产级落地
- 时间:2026-08-15
-
- PostgreSQL 16并行查询调优实战:执行计划与资源策略解析
- 时间:2026-08-15
-
- 阿里云建站产品怎么选:万小智AI建站与云企业官网区别及活动参考
- 时间:2026-08-15
-
- 年AI工具推荐精选:办公设计编程学习全场景指南
- 时间:2026-08-15
-
- 多Agent协作策略评测平台:回测过拟合检测与Walk-Forward全链路
- 时间:2026-08-15
-
- 年企业仓库管理系统选型指南与实施建议
- 时间:2026-08-15
-
- 云原生与边缘计算实战:少数民族双语考试中台重构方案
- 时间:2026-08-15
-
- RAG上线后总答非所问怎么办?黄金数据集与检索质量评测
- 时间:2026-08-15