多Agent系统架构与设计模式

🏷️ L2 📊 intermediate ⏱️ 45分钟 🏷️ 多Agent,系统架构,设计模式,Agent协调,任务分配,前沿

# 多Agent系统架构与设计模式

一、概述

为什么要学习这个主题

在当今的AI应用开发中,单一Agent的能力边界已经无法满足复杂业务场景的需求。无论是企业级智能客服、自动化运维平台,还是复杂的业务流程编排,都面临着以下挑战:

多Agent系统通过将复杂的任务拆解、分配给多个专业的Agent协同完成,已经成为解决上述问题的核心架构模式。

学完本课程能做什么

1. 设计并搭建多Agent协作系统:掌握Agent角色划分与通信机制 2. 实现任务智能分解与分配:能够设计任务调度策略 3. 构建高可用的Agent集群:掌握容错和仲裁模式 4. 优化系统性能和响应速度:理解不同设计模式的优劣 5. 解决实际业务中的多Agent协作问题:能够诊断和优化现有系统

适合人群与前置知识

---

二、核心知识点

模块1: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")

在博海学习网开始学习 →