最新下载
热门教程
- 1
- 2
- 3
- 4
- 5
- 6
- 7
- 8
- 9
- 10
Agent工程化实战:从原型到生产级多智能体系统的构建与部署
时间:2026-08-15 11:32:51 编辑:袖梨 来源:一聚教程网
Agent工程化实战:从原型到生产级多智能体系统的构建与部署需要先看清适用场景和关键步骤,避免只记结论却忽略实际限制。
Agent工程化实战:从原型到生产级多智能体系统的构建与部署
一、为什么Agent工程化如此重要?
2025年,大语言模型(LLM)的推理能力已突破实用门槛,基于LLM的智能体(Agent)从Demo走向真实业务场景。然而,原型与生产之间隔着一条巨大的工程化鸿沟:

下文会围绕上述痛点,从零构建一个生产级多智能体系统,涵盖设计、编码、测试、容器化、部署与监控全流程。
二、系统架构设计
我们设计一个智能客服 代码辅助的双Agent协作系统,包含以下角色:
Agent | 职责 | 工具 |
|---|---|---|
Planner | 分析用户问题,拆解任务,调度执行 | 任务分解、资源检索 |
Executor | 执行具体动作,如查询数据库、调用API、生成代码 | SQL执行器、代码解释器、向量检索 |
Reviewer | 校验执行结果,必要时回退或重试 | 结果校验、反思 |
整体架构如下:
代码语言:javascript复制用户请求 → 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 自定义工具
代码语言:javascript复制# 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风格提示,并加入结构化输出要求。
代码语言:javascript复制# 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、处理结果汇总。
代码语言:javascript复制# 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 _save_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._save_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._save_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任务定义
代码语言:javascript复制# 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流式)
代码语言:javascript复制# 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等。
import structloglogger = structlog.get_logger()# 在orchestrator中logger.info("agent_plan", trace_id=current_trace_id, plan=plan.dict())
指标监控:Prometheus记录请求数、延迟、LLM调用次数、工具成功率。
代码语言:javascript复制from prometheus_client import Counter, Histogram, start_http_serverREQUEST_COUNT = Counter('agent_requests_total', 'Total requests')LLM_CALL_HISTOGRAM = Histogram('llm_call_duration_seconds', 'LLM call latency')# 装饰器统计
六、容器化与Kubernetes部署
6.1 Dockerfile(多阶段构建)
代码语言:javascript复制# DockerfileFROM python:3.11-slim as builderWORKDIR /appCOPY requirements.txt .RUN pip install --no-cache-dir -r requirements.txtFROM python:3.11-slimWORKDIR /appCOPY --from=builder /usr/local/lib/python3.11/site-packages /usr/local/lib/python3.11/site-packagesCOPY . .ENV PYTHONPATH=/appEXPOSE 8000CMD ["uvicorn", "api:app", "--host", "0.0.0.0", "--port", "8000"]
6.2 Kubernetes资源编排(Deployment Service ConfigMap)
代码语言:javascript复制# deployment.yamlapiVersion: apps/v1kind: Deploymentmetadata:name: agent-apispec:replicas: 3selector:matchLabels:app: agent-apitemplate:metadata:labels:app: agent-apispec:containers:- name: apiimage: agent-system:latestports:- containerPort: 8000env:- name: REDIS_URLvalueFrom:configMapKeyRef:name: agent-configkey: redis_url- name: OPENAI_API_KEYvalueFrom:secretKeyRef:name: openai-secretkey: api_keyresources:limits:cpu: "1000m"memory: "2Gi"requests:cpu: "500m"memory: "1Gi"livenessProbe:httpGet:path: /healthport: 8000initialDelaySeconds: 30periodSeconds: 10readinessProbe:httpGet:path: /readyport: 8000---apiVersion: v1kind: Servicemetadata:name: agent-api-svcspec:selector:app: agent-apiports:- port: 80targetPort: 8000type: LoadBalancer
水平自动伸缩(HPA)基于CPU和自定义指标(如队列长度):
代码语言:javascript复制apiVersion: autoscaling/v2kind: HorizontalPodAutoscalermetadata:name: agent-api-hpaspec:scaleTargetRef:apiVersion: apps/v1kind: Deploymentname: agent-apiminReplicas: 2maxReplicas: 10metrics:- type: Resourceresource:name: cputarget:type: UtilizationaverageUtilization: 60- type: Podspods:metric:name: celery_queue_lengthtarget:type: AverageValueaverageValue: 100
七、性能优化与成本控制实践
7.1 缓存语义相似查询
使用向量数据库(如Milvus)缓存常见问答对,LLM调用前先检索,命中则直接返回,降低延迟和费用。
7.2 模型降级策略
当主模型(GPT-4o)限流或超时时,自动切换至轻量级模型(如GPT-3.5-turbo),并告警。
代码语言:javascript复制class ModelRouter:def __init__(self):self.primary = ChatOpenAI(model="gpt-4o", timeout=10)self.fallback = ChatOpenAI(model="gpt-3.5-turbo", timeout=5)async def ainvoke_with_fallback(self, prompt):try:return await self.primary.ainvoke(prompt)except Exception:logger.warning("primary model failed, using fallback")return await self.fallback.ainvoke(prompt)
7.3 流式响应与首字延迟优化
对于WebSocket,首字延迟是关键。我们采用astream_events逐token返回,减少用户等待感。
async def stream_process(self, user_input):async for event in self.llm.astream_events(user_input, version="v1"):if event["event"] == "on_chat_model_stream":yield event["data"]["chunk"].content
八、测试与质量保障
8.1 单元测试(pytest mocking)
代码语言:javascript复制# test_agent.pyimport pytestfrom unittest.mock import patchfrom orchestrator import [email protected] def test_planner_parsing():orch = AgentOrchestrator("test_sess")with patch.object(orch, '_run_executor') as mock_exec:mock_exec.return_value = "mocked result"result = await orch.process_request("查询总销售额")assert "总销售额" in result
8.2 集成测试与混沌工程
在预发布环境模拟LLM超时、Redis断连、磁盘满等故障,验证降级和恢复能力。
8.3 A/B测试与灰度发布
通过Istio实现流量切分,新版本Agent只接收5%流量,对比成功率和用户满意度。
九、总结与展望
本文从零构建了一个生产级多智能体系统,覆盖了工程化的核心维度:异步解耦、状态管理、可观测性、容器编排、性能优化。在实战中,我们深刻体会到:
Agent不是单点,而是流程:成功的关键在于清晰的编排、明确的工具契约和鲁棒的错误处理。可观测性必须内建:无法调试的Agent系统是无法运维的,OpenTelemetry 结构化日志是标配。成本与性能需动态平衡:缓存、模型路由、自适应超时等机制必不可少。未来演进方向:
自进化Agent:基于历史执行数据微调小型模型,减少对大型模型的依赖。多模态Agent:结合视觉、语音,实现更丰富的交互。联邦多智能体:跨组织协作,保护数据隐私。Agent工程化才刚刚开始,期待更多开发者将智能化能力融入业务,创造真正的价值。
相关文章
- 亿方云网页版登录入口-亿方云官网登录入口 08-15
- 朱砂痣久难消你是否能知道歌曲介绍 08-15
- POKI.免费游戏免下载入口-POKI.免费游戏入口多端同步 08-15
- 崩铁4.2银狼L5999培养材料整理一览 08-15
- picacg哔咔是无法注销吗 08-15
- 鸣潮WIKI官网入口在哪-WIKI官网入口地址分享 08-15