多 Agent 系统总乱套?这套规划优化方案让一致性从 50% 提到 98%

前言
前阵子我们做了个多 Agent 协同系统,遇到了大问题:
- Agent 之间消息乱传,不知道谁该接
- 状态对不上,你说你的我做我的
- 自主规划时冲突,互相打架
后来我花了两周优化消息路由和状态一致性,系统终于稳定了。今天聊聊工程落地的经验。
一、底层原理
1.1 多 Agent 协同的两大痛点
消息路由和状态一致性是核心难题:
graph TD
A["用户请求"] –> B["协调 Agent"]
B –> C["路由混乱"]
C –> D["不知道该发给谁"]
E["多个 Agent"] –> F["状态不一致"]
F –> G["决策冲突"]
H["自主规划"] –> I["没考虑别人在做什么"]
I –> G
核心问题:
- 没有统一的消息路由机制
- 状态没有同步
- 自主规划时缺少全局视图
1.2 常见方案对比
| 中心化协调 | 简单易实现 | 单点瓶颈 |
| 完全分布式 | 无单点 | 复杂难调试 |
| 混合架构 | 平衡 | 实现稍复杂 |
二、快速上手
先看反面教材,完全无协调:
class SimpleAgent:
def __init__(self, name):
self.name = name
self.state = {}
def send(self, msg, to):
print(f"{self.name} -> {to}: {msg}")
# 直接发,没有路由
def plan(self):
# 自己规划自己的,不管别人
return "do something"
这当然乱套。
再看改进版,加个简单协调器:
class MessageRouter:
def __init__(self):
self.agents = {}
self.message_queue = []
def register(self, agent):
self.agents[agent.name] = agent
def route(self, msg, from_agent, to_agent):
self.message_queue.append((from_agent, to_agent, msg))
if to_agent in self.agents:
self.agents[to_agent].receive(msg, from_agent)
class CoordinatedAgent:
def __init__(self, name, router):
self.name = name
self.router = router
self.local_state = {}
self.global_view = {}
def send(self, msg, to):
self.router.route(msg, self.name, to)
def receive(self, msg, from_agent):
# 更新全局视图
self.global_view[from_agent] = msg
# 处理消息
def plan(self):
# 基于全局视图规划
return "plan based on global view"
三、核心 API / 深水区
3.1 关键组件速查
| 消息路由器 | 消息分发 | 主题订阅、消息队列 |
| 状态管理器 | 状态同步 | 版本号、冲突解决 |
| 规划协调器 | 规划协同 | 资源预留、优先级 |
| 监控面板 | 可观测性 | 消息追踪、状态视图 |
3.2 消息路由实现
import uuid
from typing import Dict, List, Any
from dataclasses import dataclass
from enum import Enum
class MessageType(Enum):
REQUEST = "request"
RESPONSE = "response"
BROADCAST = "broadcast"
@dataclass
class Message:
id: str
from_agent: str
to_agent: str
type: MessageType
content: Any
timestamp: float
class TopicRouter:
def __init__(self):
self.subscribers: Dict[str, List] = {}
self.message_log: List[Message] = []
def subscribe(self, topic: str, agent):
if topic not in self.subscribers:
self.subscribers[topic] = []
self.subscribers[topic].append(agent)
def publish(self, topic: str, msg: Message):
self.message_log.append(msg)
if topic in self.subscribers:
for agent in self.subscribers[topic]:
agent.receive(msg)
class AgentWithTopics:
def __init__(self, name, router):
self.name = name
self.router = router
def send_to_topic(self, topic, content):
msg = Message(
id=str(uuid.uuid4()),
from_agent=self.name,
to_agent=topic,
type=MessageType.BROADCAST,
content=content,
timestamp=0
)
self.router.publish(topic, msg)
def receive(self, msg):
print(f"{self.name} got: {msg.content}")
3.3 状态一致性管理
class StateManager:
def __init__(self):
self.states: Dict[str, Dict] = {}
self.versions: Dict[str, int] = {}
def update_state(self, agent_name, new_state, version):
# 乐观锁
if agent_name not in self.versions or version > self.versions[agent_name]:
self.states[agent_name] = new_state
self.versions[agent_name] = version
return True
return False # 冲突
def get_global_state(self):
return {
"agents": self.states,
"versions": self.versions
}
四、实战演练
完整的多 Agent 协同框架:
import time
import uuid
from typing import Dict, List, Any, Optional
from dataclasses import dataclass, field
from enum import Enum
class AgentStatus(Enum):
IDLE = "idle"
BUSY = "busy"
ERROR = "error"
@dataclass
class AgentState:
name: str
status: AgentStatus
current_task: Optional[str] = None
version: int = 0
@dataclass
class PlanningIntent:
agent: str
action: str
resources: List[str] = field(default_factory=list)
class CoordinationService:
def __init__(self):
self.agents: Dict[str, AgentState] = {}
self.pending_intents: List[PlanningIntent] = []
self.reserved_resources: Dict[str, str] = {}
self.message_log: List[Any] = []
def register_agent(self, name: str):
self.agents[name] = AgentState(
name=name,
status=AgentStatus.IDLE,
version=0
)
def request_planning(self, intent: PlanningIntent) -> bool:
# 检查资源冲突
for res in intent.resources:
if res in self.reserved_resources and self.reserved_resources[res] != intent.agent:
return False # 资源被占用
# 预留资源
for res in intent.resources:
self.reserved_resources[res] = intent.agent
self.pending_intents.append(intent)
return True
def update_agent_state(self, name: str, status: AgentStatus, task: Optional[str]):
if name in self.agents:
state = self.agents[name]
state.status = status
state.current_task = task
state.version += 1
def get_system_view(self) -> Dict:
return {
"agents": {a.name: {"status": a.status.value, "task": a.current_task}
for a in self.agents.values()},
"reserved_resources": self.reserved_resources
}
class CollaborativeAgent:
def __init__(self, name: str, coord: CoordinationService):
self.name = name
self.coord = coord
coord.register_agent(name)
def plan(self, task: str) -> bool:
# 1. 获取全局视图
view = self.coord.get_system_view()
# 2. 生成规划意图
intent = PlanningIntent(
agent=self.name,
action=task,
resources=[f"resource_{self.name}"] # 示例资源
)
# 3. 申请规划
approved = self.coord.request_planning(intent)
if not approved:
return False
# 4. 执行规划
self.coord.update_agent_state(self.name, AgentStatus.BUSY, task)
time.sleep(0.1) # 模拟执行
self.coord.update_agent_state(self.name, AgentStatus.IDLE, None)
return True
# 使用示例
coord = CoordinationService()
agent1 = CollaborativeAgent("agent1", coord)
agent2 = CollaborativeAgent("agent2", coord)
agent1.plan("do task 1")
agent2.plan("do task 2")
print(coord.get_system_view())
五、避坑指南与最佳实践
💡 **技巧 1:资源预留规划前先预留资源,避免冲突。
⚠️ **警告 1:不要硬编码路由用主题或服务发现,不要写死 Agent 名字。
✅ **推荐:可观测性一定要有消息追踪和状态面板,不然调试到哭。
六、综合实战演示
完整的生产级方案:
import time
import json
import uuid
from typing import Dict, List, Any, Optional
from collections import defaultdict
class ObservableMessageRouter:
def __init__(self):
self.topics: Dict[str, List] = defaultdict(list)
self.trace: List[Dict] = []
def subscribe(self, topic, agent):
self.topics[topic].append(agent)
def publish(self, topic, message):
msg_id = str(uuid.uuid4())
timestamp = time.time()
self.trace.append({
"id": msg_id,
"topic": topic,
"from": message.get("from"),
"timestamp": timestamp
})
for agent in self.topics.get(topic, []):
agent.on_message(topic, message)
def get_trace(self, limit=100):
return self.trace[-limit:]
class ConsistentStateStore:
def __init__(self):
self.data: Dict[str, Any] = {}
self.versions: Dict[str, int] = {}
self.history: List[Dict] = []
def update(self, key: str, value: Any, expected_version: int) -> bool:
if key not in self.versions:
self.versions[key] = 0
if expected_version != self.versions[key]:
return False # 版本不匹配,冲突
self.versions[key] += 1
self.data[key] = value
self.history.append({
"key": key,
"value": value,
"version": self.versions[key],
"timestamp": time.time()
})
return True
def get(self, key: str):
return self.data.get(key), self.versions.get(key, 0)
class ProductionCollaborativeAgent:
def __init__(self, name: str, router: ObservableMessageRouter, store: ConsistentStateStore):
self.name = name
self.router = router
self.store = store
router.subscribe("broadcast", self)
router.subscribe(f"agent.{name}", self)
def on_message(self, topic: str, message: Dict):
print(f"[{self.name}] got message on {topic}: {message}")
def plan_action(self, action: str):
# 1. 获取最新状态
global_state, version = self.store.get("global")
# 2. 基于状态规划
plan = self._create_plan(action, global_state)
# 3. 尝试更新状态
new_state = {**(global_state or {}), "last_action": {self.name: action}}
success = self.store.update("global", new_state, version)
if success:
print(f"[{self.name}] plan executed successfully")
# 4. 广播结果
self.router.publish("broadcast", {
"from": self.name,
"action": action,
"status": "done"
})
else:
print(f"[{self.name}] conflict, retrying…")
time.sleep(0.01)
self.plan_action(action) # 重试
def _create_plan(self, action, state):
return {"action": action, "state": state}
# 生产级示例
router = ObservableMessageRouter()
store = ConsistentStateStore()
agents = [
ProductionCollaborativeAgent(f"agent{i}", router, store)
for i in range(3)
]
# 模拟协同规划
for agent in agents:
agent.plan_action(f"task_{agent.name}")
print("\\nMessage trace:")
for entry in router.get_trace():
print(entry)
七、总结
多 Agent 协同系统要做好:
- 消息路由要清晰(主题、服务发现)
- 状态同步要一致(版本号、冲突解决)
- 自主规划要协调(资源预留、全局视图)
- 系统要可观测(追踪、面板)
做好这几点,一致性就能提上去,系统就稳定了。

