首页
看点啥
插画图片
首页 看点啥 Agent工程化实战:从原型到生产级多智能体系统的构建与部署

Agent工程化实战:从原型到生产级多智能体系统的构建与部署

2026-08-15 0

Agent工程化实战:从原型到生产级多智能体系统的构建与部署的重点在于把前置条件、操作顺序和容易误判的地方分清楚。

Agent工程化实战:从原型到生产级多智能体系统的构建与部署

一、为什么Agent工程化如此重要?

2025年,大语言模型(LLM)的推理能力已突破实用门槛,基于LLM的智能体(Agent)从Demo走向真实业务场景。然而,原型与生产之间隔着一条巨大的工程化鸿沟:

Agent工程化实战:从原型到生产级多智能体系统的构建与部署

可靠性:LLM输出不确定,如何保证业务流程的稳定性?可观测性:Agent的每一次工具调用、每一步推理都可能是黑盒,如何追踪和调试?扩展性:单Agent能力有限,如何设计多智能体协作并动态伸缩?性能:LLM响应延迟高,如何异步化、缓存、流式处理?运维:模型版本更新、API限流、成本控制,如何自动化管理?

下文会围绕上述痛点,从零构建一个生产级多智能体系统,涵盖设计、编码、测试、容器化、部署与监控全流程。

二、系统架构设计

我们设计一个智能客服 代码辅助的双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调用、工具执行。

代码语言:javascript

复制

# 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_idsession_idagent_name等。

代码语言:javascript

复制

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返回,减少用户等待感。

代码语言:javascript

复制

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工程化才刚刚开始,期待更多开发者将智能化能力融入业务,创造真正的价值。

喜欢(0)

上一篇

哥谭骑士生存33天中的哥谭骑士是谁

哥谭骑士生存33天中的哥谭骑士是谁

下一篇

七界梦谭是哪个公司开发的 游戏制作公司详细解答

七界梦谭是哪个公司开发的 游戏制作公司详细解答
猜你喜欢