企业级多Agent系统构建:从Prompt Engineering到Harness Engineering实践 在实际企业级 AI 系统开发中单纯依赖 Prompt Engineering 已经难以应对复杂业务逻辑和稳定性要求。Harness Engineering 作为一种新兴的工程范式强调通过系统化的设计、编排和管控手段将多个 AI Agent 可靠地集成到生产环境中。它不仅仅是写提示词而是涵盖 Agent 生命周期管理、通信机制、状态控制、异常处理和性能监控等一系列工程实践。本文将围绕如何从零搭建一个企业级多 Agent 系统带你理解 Harness Engineering 的核心概念、技术选型、架构设计、代码实现和运维要点。1. 理解 Harness Engineering 与多 Agent 系统的关系1.1 什么是 Harness EngineeringHarness Engineering 可以理解为“缰绳工程”其核心目标是对 AI Agent 进行有效控制和调度确保它们在复杂环境中协同工作。与 Prompt Engineering 主要关注单次交互的输入输出不同Harness Engineering 更注重系统层面的可靠性、可维护性和扩展性。在企业级场景中Harness Engineering 通常涉及以下关键维度Agent 编排多个 Agent 之间的工作流定义、依赖关系和执行顺序。状态管理维护 Agent 执行过程中的上下文状态避免信息丢失或冲突。容错机制当某个 Agent 执行失败或返回异常结果时系统如何降级或重试。资源管控控制 API 调用频率、计算资源消耗和成本预算。监控观测实时追踪每个 Agent 的执行状态、输入输出和性能指标。1.2 多 Agent 系统的典型架构模式根据业务复杂度不同多 Agent 系统通常采用以下几种架构模式分层架构将 Agent 按功能层级划分如接口层、业务逻辑层、数据访问层每层由专门的 Agent 负责。流水线架构任务按固定顺序流经多个 Agent每个 Agent 完成特定处理步骤后传递给下一个。黑板架构多个 Agent 共享一个公共数据空间黑板各自监听感兴趣的数据变化并做出响应。协同架构Agent 之间通过消息传递直接通信动态协商任务分配和结果整合。在实际项目中这些架构模式常常混合使用。选择哪种模式取决于业务场景的确定性、实时性要求和系统复杂度。1.3 企业级场景的特殊要求企业级多 Agent 系统与实验性项目有显著区别必须考虑以下要求高可用性关键业务链路的 Agent 需要有备份方案和故障转移机制。安全性敏感数据的处理要符合企业安全规范避免信息泄露。审计追踪每个决策点的输入输出需要完整记录满足合规要求。性能可预测系统响应时间要在可接受范围内不能因 Agent 数量增加而急剧下降。易于集成能够与企业现有系统CRM、ERP、数据库等无缝对接。2. 搭建多 Agent 系统的技术选型与环境准备2.1 核心框架选择目前主流的多 Agent 框架主要有以下几种选择框架优势适用场景学习曲线LangGraph状态机模型清晰适合复杂工作流业务流程严格、状态转换明确的场景中等AutoGen对话式协作能力强需要多轮对话、协商决策的场景较陡CrewAI面向任务分解抽象程度高快速构建任务导向型系统平缓自建框架完全可控定制性强有特殊需求或需要深度优化的场景高对于大多数企业级项目建议从 LangGraph 开始它在状态管理和流程控制方面提供了良好的平衡。2.2 开发环境配置以 Python 为例以下是基础环境配置# 创建虚拟环境 python -m venv agent_env source agent_env/bin/activate # Linux/Mac # agent_env\Scripts\activate # Windows # 安装核心依赖 pip install langgraph langchain-openai pip install pydantic python-dotenv创建环境配置文件.env# OpenAI API 配置 OPENAI_API_KEYyour_api_key_here OPENAI_BASE_URLhttps://api.openai.com/v1 # 项目配置 AGENT_LOG_LEVELINFO MAX_CONCURRENT_AGENTS5 REQUEST_TIMEOUT302.3 项目结构设计规范的项目结构是 Harness Engineering 的基础project/ ├── agents/ # Agent 实现类 │ ├── __init__.py │ ├── base_agent.py # 基础 Agent 类 │ ├── classifier_agent.py │ └── analyzer_agent.py ├── workflows/ # 工作流定义 │ ├── __init__.py │ └── main_workflow.py ├── schemas/ # 数据模型 │ ├── __init__.py │ └── models.py ├── config/ # 配置文件 │ ├── __init__.py │ └── settings.py ├── utils/ # 工具函数 │ ├── __init__.py │ └── logger.py ├── tests/ # 测试代码 │ ├── __init__.py │ └── test_agents.py └── main.py # 应用入口这种结构确保了代码的可维护性和模块化符合企业级项目的开发规范。3. 实现基础 Agent 与工作流编排3.1 定义基础 Agent 类首先创建一个可复用的基础 Agent 类封装通用功能from abc import ABC, abstractmethod from typing import Any, Dict, Optional import logging from pydantic import BaseModel class BaseAgent(ABC): 基础 Agent 抽象类 def __init__(self, name: str, config: Dict[str, Any] None): self.name name self.config config or {} self.logger logging.getLogger(fagent.{name}) abstractmethod async def execute(self, input_data: Dict[str, Any]) - Dict[str, Any]: 执行 Agent 的核心逻辑 pass def validate_input(self, input_data: Dict[str, Any]) - bool: 输入验证 required_fields self.get_required_fields() for field in required_fields: if field not in input_data: self.logger.error(f缺少必要字段: {field}) return False return True def get_required_fields(self) - list: 获取必要输入字段 return [] def before_execute(self, input_data: Dict[str, Any]) - Dict[str, Any]: 执行前预处理 self.logger.info(fAgent {self.name} 开始执行) return input_data def after_execute(self, result: Dict[str, Any]) - Dict[str, Any]: 执行后处理 self.logger.info(fAgent {self.name} 执行完成) return result3.2 实现具体的业务 Agent以下是一个分类器 Agent 的完整实现from langchain_core.messages import HumanMessage from langchain_openai import ChatOpenAI import os from .base_agent import BaseAgent class ClassifierAgent(BaseAgent): 分类器 Agent负责对输入文本进行分类 def __init__(self): super().__init__(classifier) self.llm ChatOpenAI( modelgpt-3.5-turbo, temperature0.1, timeout30 ) self.categories [技术问题, 业务咨询, 投诉建议, 其他] def get_required_fields(self) - list: return [text] async def execute(self, input_data: Dict[str, Any]) - Dict[str, Any]: 执行分类逻辑 if not self.validate_input(input_data): return {error: 输入数据验证失败} processed_input self.before_execute(input_data) text processed_input[text] try: # 构建分类提示词 prompt f 请对以下文本进行分类只能返回分类名称 可选分类{, .join(self.categories)} 文本内容{text} 分类结果 response await self.llm.ainvoke([HumanMessage(contentprompt)]) category response.content.strip() # 验证分类结果 if category not in self.categories: category 其他 result { category: category, confidence: high, original_text: text } return self.after_execute(result) except Exception as e: self.logger.error(f分类器执行失败: {str(e)}) return {error: f分类处理异常: {str(e)}, category: 其他}3.3 使用 LangGraph 编排工作流LangGraph 通过状态机模型管理工作流以下是一个完整的工作流示例from langgraph.graph import StateGraph, END from typing import TypedDict, Annotated import operator class AgentState(TypedDict): 工作流状态定义 input_text: str category: Annotated[str, operator.add] # 分类结果 analysis: Annotated[dict, operator.add] # 分析结果 final_response: Annotated[str, operator.add] # 最终响应 errors: Annotated[list, operator.add] # 错误信息 def create_workflow(): 创建多 Agent 工作流 # 初始化图 workflow StateGraph(AgentState) # 添加节点每个节点对应一个 Agent workflow.add_node(classifier, classifier_agent_node) workflow.add_node(analyzer, analyzer_agent_node) workflow.add_node(response_generator, response_generator_node) workflow.add_node(error_handler, error_handler_node) # 设置入口点 workflow.set_entry_point(classifier) # 添加边定义执行流程 workflow.add_edge(classifier, analyzer) workflow.add_conditional_edges( analyzer, route_after_analysis, # 条件路由函数 { continue: response_generator, error: error_handler } ) workflow.add_edge(response_generator, END) workflow.add_edge(error_handler, END) return workflow.compile() def classifier_agent_node(state: AgentState) - AgentState: 分类器 Agent 节点 try: agent ClassifierAgent() result await agent.execute({text: state[input_text]}) if error in result: state[errors].append(f分类器错误: {result[error]}) else: state[category] result[category] except Exception as e: state[errors].append(f分类器异常: {str(e)}) return state def route_after_analysis(state: AgentState) - str: 分析后的路由判断 if state[errors]: return error elif state.get(analysis, {}).get(needs_human_review): return error # 需要人工审核的情况也走错误处理 else: return continue4. 关键配置与参数详解4.1 Agent 执行控制参数在多 Agent 系统中每个 Agent 的执行都需要精细控制AGENT_CONFIG { timeout: 30, # 单次执行超时时间秒 max_retries: 2, # 最大重试次数 retry_delay: 1, # 重试延迟秒 rate_limit: 10, # 每分钟最大调用次数 circuit_breaker: { # 熔断器配置 failure_threshold: 5, # 连续失败阈值 reset_timeout: 60 # 熔断恢复时间秒 } }4.2 工作流状态管理配置状态管理是 Harness Engineering 的核心关键配置包括# state_management.yaml state_config: persistence: enabled: true backend: redis # 或 database, memory ttl: 3600 # 状态保存时间秒 snapshot: enabled: true interval: 10 # 快照间隔执行步骤数 recovery: enabled: true max_replay_steps: 50 # 最大回放步数4.3 监控与日志配置企业级系统必须包含完善的监控# monitoring.py import logging from prometheus_client import Counter, Histogram # 定义指标 AGENT_EXECUTION_COUNT Counter( agent_execution_total, Agent执行总次数, [agent_name, status] ) AGENT_EXECUTION_DURATION Histogram( agent_execution_duration_seconds, Agent执行耗时, [agent_name] ) def setup_logging(): 配置结构化日志 logging.basicConfig( levellogging.INFO, format%(asctime)s - %(name)s - %(levelname)s - %(message)s, handlers[ logging.FileHandler(agent_system.log), logging.StreamHandler() ] )5. 运行验证与结果分析5.1 端到端测试流程建立完整的测试验证流程import asyncio from workflows.main_workflow import create_workflow async def test_workflow(): 测试完整工作流 # 初始化工作流 workflow create_workflow() # 测试用例 test_cases [ {input_text: 我的数据库连接总是超时怎么办}, {input_text: 我想咨询一下产品的价格信息}, {input_text: 服务态度太差了我要投诉}, {input_text: 今天天气怎么样} # 非常规问题 ] for i, test_case in enumerate(test_cases): print(f执行测试用例 {i1}: {test_case[input_text]}) try: # 执行工作流 result await workflow.ainvoke(test_case) # 验证结果 assert final_response in result, 缺少最终响应 assert result[final_response], 最终响应为空 print(f✓ 测试通过: {result[final_response][:100]}...) except Exception as e: print(f✗ 测试失败: {str(e)}) # 详细错误信息 if errors in result: for error in result[errors]: print(f 错误: {error}) # 运行测试 if __name__ __main__: asyncio.run(test_workflow())5.2 性能基准测试对企业级系统来说性能测试必不可少import time import statistics from concurrent.futures import ThreadPoolExecutor def benchmark_workflow(concurrent_users: int, requests_per_user: int): 工作流性能基准测试 async def single_user_workflow(user_id: int): durations [] workflow create_workflow() for i in range(requests_per_user): start_time time.time() try: await workflow.ainvoke({ input_text: f用户{user_id}的第{i1}个请求 }) duration time.time() - start_time durations.append(duration) except Exception as e: print(f用户{user_id}请求{i1}失败: {e}) return durations # 并发测试 async def run_concurrent_test(): tasks [] for user_id in range(concurrent_users): task asyncio.create_task(single_user_workflow(user_id)) tasks.append(task) all_durations await asyncio.gather(*tasks) # 统计结果 flat_durations [d for sublist in all_durations for d in sublist] print(f并发用户数: {concurrent_users}) print(f总请求数: {len(flat_durations)}) print(f平均响应时间: {statistics.mean(flat_durations):.2f}s) print(fP95响应时间: {statistics.quantiles(flat_durations, n20)[18]:.2f}s) print(f最大响应时间: {max(flat_durations):.2f}s) asyncio.run(run_concurrent_test()) # 执行性能测试 benchmark_workflow(concurrent_users10, requests_per_user5)6. 常见问题排查与解决方案6.1 Agent 执行失败排查当 Agent 执行出现问题时按以下顺序排查问题现象可能原因检查方式解决方案Agent 超时API 响应慢或网络问题检查超时配置和网络连接增加超时时间或添加重试机制返回结果格式错误Prompt 设计问题或模型理解偏差检查输入输出格式验证优化 Prompt 或添加结果后处理内存占用过高状态数据过大或内存泄漏检查状态序列化大小优化状态存储定期清理并发性能差资源竞争或瓶颈节点分析各节点执行时间优化慢节点增加并发控制6.2 工作流状态异常处理状态管理中的常见问题及处理class StateRecoveryManager: 状态恢复管理器 def __init__(self, workflow): self.workflow workflow self.state_backend RedisBackend() async def recover_workflow(self, workflow_id: str, target_step: int None): 恢复中断的工作流 # 获取保存的状态 saved_state await self.state_backend.get(workflow_id) if not saved_state: raise ValueError(f未找到工作流 {workflow_id} 的状态) # 检查状态完整性 if not self.validate_state(saved_state): await self.state_backend.delete(workflow_id) raise ValueError(状态数据损坏无法恢复) # 恢复到指定步骤或最新步骤 recovery_point target_step or saved_state[current_step] # 重新执行从恢复点开始的工作流 recovered_result await self.workflow.aresume( saved_state, from_steprecovery_point ) return recovered_result def validate_state(self, state: dict) - bool: 验证状态数据完整性 required_fields [workflow_id, current_step, state_data] return all(field in state for field in required_fields)6.3 监控告警配置建立有效的监控告警机制# alert_rules.yaml alert_rules: - alert: HighErrorRate expr: rate(agent_execution_total{statuserror}[5m]) 0.1 for: 2m labels: severity: warning annotations: summary: Agent 错误率过高 description: 过去5分钟内错误率超过10% - alert: WorkflowTimeout expr: agent_execution_duration_seconds 30 for: 1m labels: severity: critical annotations: summary: 工作流执行超时 description: 工作流执行时间超过30秒 - alert: APIQuotaExceeded expr: api_calls_remaining 100 labels: severity: info annotations: summary: API 配额即将用尽 description: 剩余API调用次数不足100次7. 企业级最佳实践与扩展方向7.1 安全与合规实践在企业环境中安全是首要考虑class SecurityManager: 安全管理器 def __init__(self): self.sensitive_patterns [ r\b\d{4}[- ]?\d{4}[- ]?\d{4}[- ]?\d{4}\b, # 信用卡号 r\b\d{3}[- ]?\d{2}[- ]?\d{4}\b, # 社会安全号 # 添加更多敏感数据模式 ] def sanitize_input(self, text: str) - str: 输入数据脱敏 for pattern in self.sensitive_patterns: text re.sub(pattern, [REDACTED], text) return text def audit_log(self, agent_name: str, input_data: dict, output_data: dict): 审计日志记录 audit_entry { timestamp: datetime.utcnow().isoformat(), agent: agent_name, input_hash: hashlib.sha256(str(input_data).encode()).hexdigest(), output_hash: hashlib.sha256(str(output_data).encode()).hexdigest(), user_id: self.get_current_user_id() } # 写入安全审计日志 self.audit_logger.info(json.dumps(audit_entry))7.2 性能优化策略随着系统规模扩大性能优化变得重要缓存策略对频繁使用的分类结果建立缓存对稳定的业务规则预计算使用 Redis 或 Memcached 作为缓存后端异步处理非实时任务使用消息队列异步处理批量处理相似请求减少 API 调用实现请求合并和批量发送资源池化建立 Agent 实例池避免重复初始化连接池管理数据库和 API 连接内存池优化大对象分配7.3 扩展性设计为未来扩展做好准备class PluginManager: 插件管理器支持动态扩展 def __init__(self): self.plugins {} self.hooks { pre_agent_execute: [], post_agent_execute: [], workflow_completed: [] } def register_plugin(self, name: str, plugin_class): 注册插件 self.plugins[name] plugin_class def add_hook(self, hook_point: str, callback): 添加钩子函数 self.hooks[hook_point].append(callback) async def execute_hooks(self, hook_point: str, context: dict): 执行钩子 for callback in self.hooks.get(hook_point, []): try: await callback(context) except Exception as e: logging.error(f钩子执行失败 {hook_point}: {e}) # 使用示例 plugin_manager PluginManager() # 注册一个数据验证插件 plugin_manager.register_plugin(data_validator, DataValidatorPlugin) # 添加预处理钩子 plugin_manager.add_hook(pre_agent_execute, lambda ctx: validate_input_data(ctx))7.4 部署与运维建议生产环境部署需要考虑容器化部署FROM python:3.11-slim WORKDIR /app COPY requirements.txt . RUN pip install -r requirements.txt COPY . . EXPOSE 8000 CMD [python, main.py]配置管理使用环境变量区分配置敏感信息通过密钥管理服务获取配置变更需要滚动重启健康检查# health_check.py async def health_check(): 健康检查端点 checks { database: check_database_connection(), redis: check_redis_connection(), external_api: check_api_availability() } all_healthy all(checks.values()) status_code 200 if all_healthy else 503 return { status: healthy if all_healthy else unhealthy, checks: checks, timestamp: datetime.utcnow().isoformat() }Harness Engineering 的真正价值在于将 AI Agent 从实验性技术转变为可靠的生产力工具。通过系统化的工程实践企业可以构建出既智能又稳定的多 Agent 系统在保证业务连续性的同时享受 AI 技术带来的效率提升。实际项目中建议从小规模试点开始逐步验证每个环节的可靠性再扩展到更复杂的业务场景。