多Agent协同实战:我是如何设计智能协作系统的
·
多Agent协同实战:我是如何设计智能协作系统的

前言
最近在做一个复杂的AI工作流系统,需要多个Agent协同完成任务。
一开始简单地让Agent直接通信,结果出现了状态不一致的问题。
后来引入了消息路由和状态管理机制,系统稳定性提升了80%。
这篇文章分享我的设计经验。
一、底层原理
1.1 核心机制
多Agent协同的关键是消息路由和状态一致性:
graph TD
A[任务输入] --> B[任务分发器]
B --> C[Agent 1]
B --> D[Agent 2]
B --> E[Agent 3]
C --> F[消息总线]
D --> F
E --> F
F --> G[状态管理器]
G --> H{状态一致?}
H -->|是| I[任务完成]
H -->|否| J[冲突解决]
J --> F
关键组件:
| 组件 | 功能 | 作用 |
|---|---|---|
| 消息总线 | 消息路由 | 统一通信 |
| 状态管理器 | 状态同步 | 一致性保障 |
| 任务分发器 | 任务分配 | 负载均衡 |
| 冲突解决器 | 冲突处理 | 协调矛盾 |
1.2 与同类方案的对比
| 架构 | 扩展性 | 一致性 | 复杂度 |
|---|---|---|---|
| 中心化 | 低 | 高 | 低 |
| 去中心化 | 高 | 中 | 高 |
| 混合架构 | 高 | 高 | 中 |
二、快速上手
from langchain.agents import AgentExecutor
from langchain.llms import OpenAI
# 创建多个Agent
agent1 = create_agent("分析专家")
agent2 = create_agent("执行专家")
agent3 = create_agent("总结专家")
# 简单的协同流程
def simple_workflow(task):
# Agent1分析任务
analysis = agent1.run(task)
# Agent2执行任务
result = agent2.run(analysis)
# Agent3总结结果
summary = agent3.run(result)
return summary
三、核心 API / 深水区
3.1 核心方法速查
| 方法 | 功能 | 适用场景 |
|---|---|---|
AgentExecutor() | Agent执行器 | 单Agent执行 |
ChatOpenAI() | 对话模型 | 多轮对话 |
Memory() | 记忆管理 | 状态保持 |
EventEmitter() | 事件发布 | 消息通知 |
StateMachine() | 状态机 | 流程控制 |
3.2 生产级配置
from langchain.chains import LLMChain
from langchain.prompts import PromptTemplate
from langchain.chat_models import ChatOpenAI
class Coordinator:
def __init__(self):
self.agents = {}
self.state = {}
self.message_bus = MessageBus()
def register_agent(self, name, agent):
self.agents[name] = agent
def dispatch(self, task):
# 分析任务类型
analysis = self.analyze_task(task)
# 分配给合适的Agent
agent_name = self.select_agent(analysis)
# 执行任务
result = self.agents[agent_name].run(task)
# 更新状态
self.update_state(agent_name, result)
return result
def analyze_task(self, task):
prompt = PromptTemplate(
input_variables=["task"],
template="分析任务类型:{task}"
)
chain = LLMChain(llm=ChatOpenAI(), prompt=prompt)
return chain.run(task)
3.3 高级定制
# 消息总线实现
class MessageBus:
def __init__(self):
self.subscribers = {}
def subscribe(self, topic, callback):
if topic not in self.subscribers:
self.subscribers[topic] = []
self.subscribers[topic].append(callback)
def publish(self, topic, message):
if topic in self.subscribers:
for callback in self.subscribers[topic]:
callback(message)
四、实战演练
场景:多Agent代码审查流程
def code_review_workflow(code):
# 第一步:语法检查
syntax_agent = AgentExecutor.from_agent_and_tools(
agent=syntax_agent,
tools=[check_syntax]
)
syntax_result = syntax_agent.run(code)
# 第二步:安全检查
security_agent = AgentExecutor.from_agent_and_tools(
agent=security_agent,
tools=[check_security]
)
security_result = security_agent.run(code)
# 第三步:性能检查
performance_agent = AgentExecutor.from_agent_and_tools(
agent=performance_agent,
tools=[check_performance]
)
performance_result = performance_agent.run(code)
# 汇总结果
summary = f"""
语法检查:{syntax_result}
安全检查:{security_result}
性能检查:{performance_result}
"""
return summary
五、避坑指南与最佳实践
💡 技巧:使用状态机管理流程
from transitions import Machine
class WorkflowStateMachine:
states = ['pending', 'analyzing', 'executing', 'reviewing', 'completed']
def __init__(self):
self.machine = Machine(
model=self,
states=WorkflowStateMachine.states,
initial='pending'
)
# 定义状态转换
self.machine.add_transition(
trigger='start',
source='pending',
dest='analyzing'
)
self.machine.add_transition(
trigger='analyze_done',
source='analyzing',
dest='executing'
)
⚠️ 警告:避免状态不一致
# 错误示例:无状态同步
def badWorkflow(task):
agent1_result = agent1.run(task)
agent2_result = agent2.run(task) # 不知道agent1的结果
# 可能导致冲突
return combine(agent1_result, agent2_result)
# 正确做法:共享状态
def goodWorkflow(task):
shared_state = {}
agent1_result = agent1.run(task, state=shared_state)
shared_state['agent1_result'] = agent1_result
agent2_result = agent2.run(task, state=shared_state) # 知道agent1的结果
return combine(agent1_result, agent2_result)
✅ 推荐:使用事件驱动架构
class EventDrivenCoordinator:
def __init__(self):
self.message_bus = MessageBus()
self.message_bus.subscribe('task.completed', self.on_task_completed)
def on_task_completed(self, message):
task_id = message['task_id']
result = message['result']
# 触发下一步
self.trigger_next(task_id, result)
六、综合实战演示
from langchain.agents import initialize_agent, AgentType
from langchain.tools import Tool
from langchain.chat_models import ChatOpenAI
from collections import defaultdict
class MultiAgentSystem:
def __init__(self):
self.llm = ChatOpenAI(model_name="gpt-4")
self.agents = {}
self.state = defaultdict(dict)
self.message_bus = MessageBus()
# 注册消息处理器
self.message_bus.subscribe("agent.completed", self.handle_completion)
def create_agent(self, name, tools, description):
agent = initialize_agent(
tools,
self.llm,
agent=AgentType.STRUCTURED_CHAT_ZERO_SHOT_REACT_DESCRIPTION,
verbose=True
)
self.agents[name] = {
'agent': agent,
'description': description
}
return agent
def dispatch(self, task, context):
# 分析任务选择Agent
agent_name = self.select_agent(task)
# 执行任务
result = self.agents[agent_name]['agent'].run(task)
# 发布完成事件
self.message_bus.publish("agent.completed", {
'agent': agent_name,
'task': task,
'result': result,
'context': context
})
return result
def select_agent(self, task):
prompt = f"""
任务:{task}
可用Agent:{[(name, info['description']) for name, info in self.agents.items()]}
请选择最合适的Agent:
"""
result = self.llm.predict(prompt)
return result.strip()
def handle_completion(self, message):
agent_name = message['agent']
result = message['result']
context = message['context']
# 更新状态
self.state[context['task_id']][agent_name] = result
# 检查是否完成
if self.is_complete(context['task_id']):
self.finalize(context['task_id'])
def is_complete(self, task_id):
# 检查所有必要Agent是否完成
required_agents = ['analyzer', 'executor', 'reviewer']
state = self.state[task_id]
for agent in required_agents:
if agent not in state:
return False
return True
def finalize(self, task_id):
state = self.state[task_id]
summary = f"""
任务完成总结:
分析结果:{state['analyzer']}
执行结果:{state['executor']}
审查结果:{state['reviewer']}
"""
print(summary)
# 使用示例
system = MultiAgentSystem()
# 创建Agent
system.create_agent(
name='analyzer',
tools=[],
description='任务分析专家'
)
system.create_agent(
name='executor',
tools=[],
description='任务执行专家'
)
system.create_agent(
name='reviewer',
tools=[],
description='结果审查专家'
)
# 提交任务
system.dispatch(
task='分析并执行代码优化任务',
context={'task_id': 'task_001'}
)
七、总结
多Agent协同的核心是消息路由和状态管理。
关键要点:
- 使用消息总线统一通信
- 维护全局状态保证一致性
- 使用状态机管理流程
- 实现事件驱动架构
核心收获: 好的架构设计是多Agent协同成功的关键。
更多推荐



所有评论(0)