多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协同的核心是消息路由和状态管理。

关键要点:

  1. 使用消息总线统一通信
  2. 维护全局状态保证一致性
  3. 使用状态机管理流程
  4. 实现事件驱动架构

核心收获: 好的架构设计是多Agent协同成功的关键。

Logo

这里是“一人公司”的成长家园。我们提供从产品曝光、技术变现到法律财税的全栈内容,并连接云服务、办公空间等稀缺资源,助你专注创造,无忧运营。

更多推荐