# 多Agent系统架构与设计模式
在当今的AI应用开发中,单一Agent的能力边界已经无法满足复杂业务场景的需求。无论是企业级智能客服、自动化运维平台,还是复杂的业务流程编排,都面临着以下挑战:
多Agent系统通过将复杂的任务拆解、分配给多个专业的Agent协同完成,已经成为解决上述问题的核心架构模式。
1. 设计并搭建多Agent协作系统:掌握Agent角色划分与通信机制 2. 实现任务智能分解与分配:能够设计任务调度策略 3. 构建高可用的Agent集群:掌握容错和仲裁模式 4. 优化系统性能和响应速度:理解不同设计模式的优劣 5. 解决实际业务中的多Agent协作问题:能够诊断和优化现有系统
---
在多Agent系统中,每个Agent都有明确的角色定义。常见的角色包括:
| 角色类型 | 职责描述 | 典型应用 | |---------|---------|---------| | Orchestrator(协调者) | 任务分发、结果汇总、异常处理 | 主控Agent | | Specialist(专家) | 专注于特定领域的任务执行 | 代码生成、数据分析 | | Validator(验证者) | 检查输出质量、进行交叉验证 | 质量保证 | | Memory(记忆体) | 维护上下文和长期记忆 | 对话历史管理 |
from typing import List, Dict, Any
from enum import Enum
class AgentRole(Enum):
ORCHESTRATOR = "orchestrator"
SPECIALIST = "specialist"
VALIDATOR = "validator"
MEMORY = "memory"
class BaseAgent:
def __init__(self, name: str, role: AgentRole):
self.name = name
self.role = role
self.context = {}
def process(self, task: Dict[str, Any]) -> Dict[str, Any]:
raise NotImplementedError
class OrchestratorAgent(BaseAgent):
def __init__(self, name: str):
super().__init__(name, AgentRole.ORCHESTRATOR)
self.specialists = []
def register_specialist(self, agent: BaseAgent):
self.specialists.append(agent)
def process(self, task: Dict[str, Any]) -> Dict[str, Any]:
# 任务分解逻辑
subtasks = self._decompose_task(task)
results = []
for subtask in subtasks:
specialist = self._select_specialist(subtask)
result = specialist.process(subtask)
results.append(result)
return self._aggregate_results(results)
def _decompose_task(self, task):
# 任务分解实现
pass
def _select_specialist(self, subtask):
# 选择最合适的专家Agent
pass
def _aggregate_results(self, results):
# 汇总结果
pass
模块2:通信与协调机制
原理讲解
Agent间的通信机制是系统设计的核心。主要模式包括:
1. 消息队列模式:异步通信,解耦Agent
2. 事件驱动模式:基于事件的触发机制
3. RPC模式:同步调用,实时性高
4. 黑板模式:共享内存,适合协作推理
通信模式对比
| 模式 | 优点 | 缺点 | 适用场景 |
|-----|------|------|---------|
| 消息队列 | 高吞吐、解耦 | 延迟较高 | 非实时任务 |
| 事件驱动 | 灵活、可扩展 | 调试困难 | 复杂流程 |
| RPC | 低延迟、简单 | 耦合度高 | 实时交互 |
| 黑板模式 | 共享状态、协作好 | 并发控制复杂 | 推理系统 |
代码示例
import asyncio
from typing import Dict, Any, Callable
from dataclasses import dataclass
from enum import Enum
class MessageType(Enum):
TASK_ASSIGN = "task_assign"
TASK_RESULT = "task_result"
STATUS_UPDATE = "status_update"
ERROR = "error"
@dataclass
class Message:
sender: str
receiver: str
message_type: MessageType
payload: Dict[str, Any]
correlation_id: str = None
class MessageBroker:
def __init__(self):
self.queues: Dict[str, asyncio.Queue] = {}
self.subscribers: Dict[str, Callable] = {}
def register_agent(self, agent_id: str):
self.queues[agent_id] = asyncio.Queue()
def subscribe(self, event_type: str, callback: Callable):
if event_type not in self.subscribers:
self.subscribers[event_type] = []
self.subscribers[event_type].append(callback)
async def send_message(self, message: Message):
if message.receiver in self.queues:
await self.queues[message.receiver].put(message)
# 触发事件订阅
if message.message_type.value in self.subscribers:
for callback in self.subscribers[message.message_type.value]:
await callback(message)
async def receive_message(self, agent_id: str) -> Message:
if agent_id in self.queues:
return await self.queues[agent_id].get()
return None
# 使用示例
async def main():
broker = MessageBroker()
broker.register_agent("agent_1")
broker.register_agent("agent_2")
# 发送任务
await broker.send_message(Message(
sender="orchestrator",
receiver="agent_1",
message_type=MessageType.TASK_ASSIGN,
payload={"task": "analyze_data"}
))
# 接收任务
message = await broker.receive_message("agent_1")
print(f"Received: {message.payload}")
模块3:任务分解与分配策略
原理讲解
任务分解是多Agent系统的核心能力。主要策略包括:
1. 层次分解:将大任务递归分解为子任务
2. 功能分解:按功能维度拆分
3. 数据分解:按数据分区拆分
4. 时序分解:按时间顺序拆分
分配算法对比
| 算法 | 复杂度 | 适用场景 | 示例 |
|-----|--------|---------|------|
| 轮询 | O(1) | 同构Agent | 负载均衡 |
| 最短队列 | O(n) | 异构Agent | 任务调度 |
| 能力匹配 | O(n*m) | 专业Agent | 精准分配 |
| 拍卖算法 | O(n^2) | 复杂任务 | 资源优化 |
代码示例
import heapq
from typing import List, Dict, Any
class TaskDecomposer:
def __init__(self):
self.decomposition_rules = {}
def register_rule(self, task_type: str, rule_func):
self.decomposition_rules[task_type] = rule_func
def decompose(self, task: Dict[str, Any]) -> List[Dict[str, Any]]:
task_type = task.get("type", "default")
if task_type in self.decomposition_rules:
return self.decomposition_rules[task_type](task)
return [task] # 不可分解
class TaskScheduler:
def __init__(self):
self.agent_capabilities = {}
self.task_queue = []
self.agent_load = {}
def register_agent(self, agent_id: str, capabilities: List[str], max_load: int = 5):
self.agent_capabilities[agent_id] = capabilities
self.agent_load[agent_id] = 0
def assign_task(self, task: Dict[str, Any]) -> str:
required_capability = task.get("required_capability")
available_agents = [
agent_id for agent_id, caps in self.agent_capabilities.items()
if required_capability in caps and self.agent_load[agent_id] < 5
]
if not available_agents:
# 使用最短队列策略
return min(self.agent_load, key=self.agent_load.get)
# 选择负载最小的可用Agent
selected_agent = min(available_agents, key=lambda x: self.agent_load[x])
self.agent_load[selected_agent] += 1
return selected_agent
def complete_task(self, agent_id: str):
if agent_id in self.agent_load:
self.agent_load[agent_id] = max(0, self.agent_load[agent_id] - 1)
# 使用示例
scheduler = TaskScheduler()
scheduler.register_agent("data_agent", ["data_analysis"])
scheduler.register_agent("code_agent", ["code_generation"])
scheduler.register_agent("qa_agent", ["validation"])
task = {"type": "analysis", "required_capability": "data_analysis"}
assigned_agent = scheduler.assign_task(task)
print(f"Task assigned to: {assigned_agent}")
模块4:容错与仲裁模式
原理讲解
多Agent系统必须处理各种故障场景。关键模式包括:
1. 主从模式:主Agent故障时从Agent接管
2. 仲裁模式:多个Agent投票决策
3. 检查点模式:定期保存状态,故障时恢复
4. 熔断模式:检测到故障时自动降级
仲裁策略对比
| 策略 | 容错性 | 性能开销 | 适用场景 |
|-----|--------|---------|---------|
| 简单多数 | 高 | 低 | 快速决策 |
| 加权投票 | 中 | 中 | 信任度不同 |
| 共识算法 | 极高 | 高 | 关键决策 |
| 降级策略 | 中 | 低 | 非关键任务 |
代码示例
import asyncio
from typing import List, Dict, Any
from collections import Counter
class VotingArbitrator:
def __init__(self, agents: List[str]):
self.agents = agents
self.agent_weights = {agent: 1.0 for agent in agents}
self.failure_threshold = 3
def set_agent_weight(self, agent_id: str, weight: float):
self.agent_weights[agent_id] = weight
async def arbitrate(self, task: Dict[str, Any], votes: List[Dict[str, Any]]) -> Dict[str, Any]:
# 收集投票结果
vote_counts = Counter()
vote_weighted = {}
for vote in votes:
agent_id = vote["agent_id"]
decision = vote["decision"]
weight = self.agent_weights.get(agent_id, 1.0)
vote_counts[decision] += 1
vote_weighted[decision] = vote_weighted.get(decision, 0) + weight
# 简单多数决策
if len(vote_counts) > 0:
majority_decision = vote_counts.most_common(1)[0][0]
return {
"decision": majority_decision,
"confidence": vote_counts[majority_decision] / len(votes),
"weighted_confidence": vote_weighted.get(majority_decision, 0) / sum(vote_weighted.values())
}
return {"decision": "abstain", "confidence": 0.0}
class CircuitBreaker:
def __init__(self, threshold: int = 5, recovery_time: int = 30):
self.failure_count = 0
self.threshold = threshold
self.recovery_time = recovery_time
self.state = "closed" # closed, open, half-open
self.last_failure_time = 0
async def call(self, func, *args, **kwargs):
if self.state == "open":
current_time = asyncio.get_event_loop().time()
if current_time - self.last_failure_time > self.recovery_time:
self.state = "half-open"
else:
raise Exception("Circuit breaker is open")
try:
result = await func(*args, **kwargs)
if self.state == "half-open":
self.state = "closed"
self.failure_count = 0
return result
except Exception as e:
self.failure_count += 1
self.last_failure_time = asyncio.get_event_loop().time()
if self.failure_count >= self.threshold:
self.state = "open"
raise e
# 使用示例
async def main():
circuit_breaker = CircuitBreaker(threshold=3, recovery_time=10)
async def risky_operation():
# 可能失败的操作
raise Exception("Operation failed")
try:
result = await circuit_breaker.call(risky_operation)
except Exception as e:
print(f"Operation failed: {e}")
print(f"Circuit state: {circuit_breaker.state}")
---
三、实操步骤
步骤1:搭建基础多Agent框架
创建一个简单的多Agent系统,包含协调者和两个专家Agent:
# agent_framework.py
import asyncio
from typing import Dict, Any
class SimpleAgentSystem:
def __init__(self):
self.agents = {}
self.task_queue = asyncio.Queue()
def add_agent(self, name: str, agent_func):
self.agents[name] = agent_func
async def process_task(self, task: Dict[str, Any]):
# 任务分解
subtasks = self._decompose(task)
# 并行执行
tasks = []
for subtask in subtasks:
agent_name = subtask["agent"]
if agent_name in self.agents:
tasks.append(self.agents[agent_name](subtask))
results = await asyncio.gather(*tasks, return_exceptions=True)
return self._aggregate(results)
def _decompose(self, task):
# 简单分解示例
return [
{"agent": "analyzer", "data": task.get("data", "")},
{"agent": "formatter", "data": task.get("data", "")}
]
def _aggregate(self, results):
return {"status": "completed", "results": results}
# 运行示例
async def main():
system = SimpleAgentSystem()
async def analyzer(task):
await asyncio.sleep(1)
return f"Analyzed: {task['data']}"
async def formatter(task):
await asyncio.sleep(0.5)
return f"Formatted: {task['data']}"
system.add_agent("analyzer", analyzer)
system.add_agent("formatter", formatter)
result = await system.process_task({"data": "test_data"})
print(f"Final result: {result}")
asyncio.run(main())
预期效果:系统会并行执行分析和格式化任务,最终输出汇总结果。
步骤2:实现消息通信机制
# message_system.py
import asyncio
from collections import defaultdict
class MessageSystem:
def __init__(self):
self.agents = {}
self.mailboxes = defaultdict(asyncio.Queue)
def register_agent(self, agent_id: str):
self.agents[agent_id] = agent_id
async def send(self, sender: str, receiver: str, message: Dict[str, Any]):
if receiver in self.agents:
await self.mailboxes[receiver].put({
"sender": sender,
"content": message,
"timestamp": asyncio.get_event_loop().time()
})
async def receive(self, agent_id: str) -> Dict[str, Any]:
if agent_id in self.agents:
return await self.mailboxes[agent_id].get()
return None
# 测试消息系统
async def test_messaging():
ms = MessageSystem()
ms.register_agent("coordinator")
ms.register_agent("worker_1")
# 发送消息
await ms.send("coordinator", "worker_1", {"task": "process_data"})
# 接收消息
message = await ms.receive("worker_1")
print(f"Worker received: {message}")
asyncio.run(test_messaging())
预期效果:Agent间能够通过异步消息队列进行通信。
步骤3:实现故障恢复机制
# fault_tolerance.py
import asyncio
import random
class FaultTolerantAgent:
def __init__(self, max_retries=3, timeout=5):
self.max_retries = max_retries
self.timeout = timeout
self.health_status = "healthy"
async def execute_with_retry(self, func, *args, **kwargs):
for attempt in range(self.max_retries):
try:
result = await asyncio.wait_for(
func(*args, **kwargs),
timeout=self.timeout
)
self.health_status = "healthy"
return result
except asyncio.TimeoutError:
print(f"Attempt {attempt + 1} timed out")