多智能体互联时的循环协调:GB/Z 185框架下的Multi-Agent Loop Engineering(Python 3.10+ 实现 + 可运行代码)
多智能体互联时的循环协调:GB/Z 185框架下的Multi-Agent Loop Engineering(Python 3.10+ 实现 + 可运行代码)
一句话总结:同步看事件总线,冲突看乐观锁,恢复看Checkpoint。三种协调模式不是三选一的选择题,而是可以嵌套叠加的组合拳——这是我把三种模式逐一实现、逐个跑通后最大的体会。
📅 本文版本:2026年7月 · Python 3.10+ · asyncio(Python内置库,无需额外安装)
📌 前置阅读:《GB/Z 185合规的Agent Loop设计》 · 《A2A协议实战》 · 《Loop Engineering四层架构》
导读:我自己就踩过一次:去年做一个多Agent文档审核系统,三个Agent各负责一轮——一个查条款、一个写报告、一个做终审。我把每个Agent的Loop都写到自认为完美了,结果联调那天翻车了:查条款的Agent还在循环里没出来,写报告的Agent拿不到数据直接空转输出了一篇空白报告,终审Agent发了疯一样重试了12次。盯着日志看了两小时,才搞明白——问题不在任何一个Agent的Loop逻辑,在于它们的循环根本没有对齐节奏。
文章目录
一、前置条件与版本说明
⚠️ 本文环境:Python 3.10+(推荐3.11+),asyncio标准库(Python内置,无需pip安装)
⚠️ 适用范围:本文代码基于Python asyncio实现,适用于所有支持async/await的Python 3.7+环境
⚠️ 生产环境提示:本文Demo代码是用来教学的,上生产还缺三样:持久化存储(数据不丢)、分布式锁(不打架)、日志监控(出事了能查)
# 验证Python版本(需要3.7+,推荐3.11+)
python --version
# 预期输出:Python 3.11.x 或 3.10.x
# 验证asyncio可用性
import asyncio
import sys
print(f"Python {sys.version}")
print(f"asyncio module: {asyncio.__name__}") # 预期输出:asyncio module: asyncio
二、为什么多Agent的循环需要专门协调?
2.1 单Agent Loop vs Multi-Agent Loop:问题质变
单个Agent的Loop Engineering只需要关心:
- 迭代什么时候停止?
- 输出质量怎么验证?
- 异常了怎么回滚?
但当你有多个Agent时,问题变成了:
- Agent A完成了,Agent B怎么知道?(同步问题)
- 两个Agent同时修改同一份数据,听谁的?(冲突问题)
- Agent B卡住了,Agent A该等多久?(超时问题)
- Agent A重启了,怎么恢复到之前的状态?(状态一致性问题)
核心洞察:Multi-Agent不是"多个单Agent的叠加",而是一个分布式系统。每个Agent有自己的循环节奏,循环之间需要像操作系统调度进程一样被协调。
2.2 GB/Z 185 8.1条款:多智能体协作的规范要求
GB/Z 185标准8.1条款定义了多智能体协作的核心要求:
| 标准条款 | 要求内容 | 对循环协调的意义 | 本文实现 |
|---|---|---|---|
| 8.1.1 协作协议 | 定义Agent间的消息格式和交互规则 | 决定了循环间的通信方式 | A2A协议作为通信层 |
| 8.1.2 角色分工 | 明确各Agent的职责(发起者/执行者/仲裁者) | 决定了循环协调的拓扑结构 | 主从/对等/管道三种模式 |
| 8.1.3 冲突解决 | 定义多Agent冲突时的仲裁机制 | 决定了循环冲突的处理策略 | 乐观锁/悲观锁/仲裁者模式 |
| 8.1.4 状态同步 | 确保各Agent对共享状态的一致性认知 | 决定了循环状态的同步机制 | 状态机+事件驱动 |
三、三种Multi-Agent循环协调模式
根据GB/Z 185 8.1.2的角色分工要求,多Agent循环协调有三种典型拓扑:
3.1 模式一:主从协调(Master-Worker)
架构:一个调度Agent(Master)控制多个执行Agent(Worker)的循环节奏。
适用场景:
- 任务可以被明确拆解为子任务
- 需要统一调度、统一汇总
- 各Worker之间无依赖关系
核心问题:Master怎么知道Worker的Loop是否完成?
Python实现:
import asyncio
import inspect
from typing import List, Dict, Callable
from dataclasses import dataclass
from enum import Enum
class WorkerStatus(Enum):
"""Worker Agent循环状态。"""
IDLE = "idle"
RUNNING = "running"
COMPLETED = "completed"
FAILED = "failed"
TIMEOUT = "timeout"
@dataclass
class WorkerResult:
"""Worker执行结果。"""
worker_id: str
status: WorkerStatus
output: str
iterations: int
duration_ms: float
class MasterWorkerCoordinator:
"""
主从协调器:Master Loop控制多个Worker Loop的执行节奏。
GB/Z 185 8.1.2 角色分工:Master=发起者/仲裁者,Worker=执行者
"""
def __init__(self, max_workers: int = 5, default_timeout: float = 30.0):
self.max_workers = max_workers
self.default_timeout = default_timeout
self.workers: Dict[str, Callable] = {} # worker_id -> worker_func
self.worker_status: Dict[str, WorkerStatus] = {}
self.worker_results: Dict[str, WorkerResult] = {}
def register_worker(self, worker_id: str, worker_func: Callable):
"""注册一个Worker Agent。"""
self.workers[worker_id] = worker_func
self.worker_status[worker_id] = WorkerStatus.IDLE
async def execute_task(self, master_task: str,
worker_assignments: Dict[str, str]) -> Dict[str, WorkerResult]:
"""
Master分配任务给多个Worker,等待全部完成后汇总。
Args:
master_task: Master的原始任务描述
worker_assignments: {worker_id: sub_task_description}
"""
print(f"[Master] 收到任务: {master_task}")
print(f"[Master] 拆解为 {len(worker_assignments)} 个子任务")
# 1. Master循环:创建所有Worker任务
tasks = []
for worker_id, sub_task in worker_assignments.items():
if worker_id not in self.workers:
print(f"[Master] ❌ Worker '{worker_id}' 未注册")
continue
self.worker_status[worker_id] = WorkerStatus.RUNNING
# 异步启动Worker Loop
task = asyncio.create_task(
self._run_worker_loop(worker_id, sub_task),
name=worker_id
)
tasks.append(task)
# 2. Master循环:等待所有Worker完成(或超时)
print(f"[Master] 等待 {len(tasks)} 个Worker完成...")
# 使用asyncio.gather等待,但带超时控制
try:
# 等待所有任务完成,或整体超时
await asyncio.wait_for(
asyncio.gather(*tasks, return_exceptions=True),
timeout=self.default_timeout
)
except asyncio.TimeoutError:
print(f"[Master] ⚠️ 整体任务超时(>{self.default_timeout}s)")
self._handle_timeout()
# 3. Master循环:汇总结果
return self._aggregate_results()
async def _run_worker_loop(self, worker_id: str, sub_task: str):
"""
执行单个Worker的Loop。
这是Worker自己的循环,被Master异步调度。
"""
import time
start = time.time()
try:
print(f" [{worker_id}] 开始执行子任务: {sub_task}")
# 调用Worker的Agent Loop
worker_func = self.workers[worker_id]
result = await worker_func(sub_task) if inspect.iscoroutinefunction(worker_func) else worker_func(sub_task)
duration = (time.time() - start) * 1000
self.worker_results[worker_id] = WorkerResult(
worker_id=worker_id,
status=WorkerStatus.COMPLETED,
output=str(result),
iterations=1, # 简化:实际应从Worker Loop获取
duration_ms=duration
)
self.worker_status[worker_id] = WorkerStatus.COMPLETED
print(f" [{worker_id}] ✅ 完成 ({duration:.0f}ms)")
except Exception as e:
duration = (time.time() - start) * 1000
self.worker_results[worker_id] = WorkerResult(
worker_id=worker_id,
status=WorkerStatus.FAILED,
output=str(e),
iterations=0,
duration_ms=duration
)
self.worker_status[worker_id] = WorkerStatus.FAILED
print(f" [{worker_id}] ❌ 失败: {e}")
def _handle_timeout(self):
"""处理超时:标记未完成的Worker为TIMEOUT。"""
for worker_id, status in self.worker_status.items():
if status == WorkerStatus.RUNNING:
self.worker_status[worker_id] = WorkerStatus.TIMEOUT
self.worker_results[worker_id] = WorkerResult(
worker_id=worker_id,
status=WorkerStatus.TIMEOUT,
output="执行超时",
iterations=0,
duration_ms=self.default_timeout * 1000
)
def _aggregate_results(self) -> Dict[str, WorkerResult]:
"""Master汇总所有Worker结果。"""
completed = sum(1 for r in self.worker_results.values() if r.status == WorkerStatus.COMPLETED)
failed = sum(1 for r in self.worker_results.values() if r.status in [WorkerStatus.FAILED, WorkerStatus.TIMEOUT])
print(f"\n[Master] 任务汇总:")
print(f" 成功: {completed}/{len(self.worker_results)}")
print(f" 失败/超时: {failed}/{len(self.worker_results)}")
return self.worker_results
# ========== 使用示例:主从协调 ==========
async def demo_master_worker():
"""演示:调度Agent分配任务给数据Agent和文档Agent。"""
coordinator = MasterWorkerCoordinator(default_timeout=10.0)
# 注册Worker
async def data_agent_task(task: str):
await asyncio.sleep(0.5) # 模拟数据查询
return f"数据查询结果: Q2销售额1200万"
async def doc_agent_task(task: str):
await asyncio.sleep(0.8) # 模拟文档生成
return f"文档生成结果: Q2销售报告.docx"
coordinator.register_worker("data_agent", data_agent_task)
coordinator.register_worker("doc_agent", doc_agent_task)
# Master分配任务
results = await coordinator.execute_task(
master_task="生成Q2销售报告",
worker_assignments={
"data_agent": "查询Q2销售数据",
"doc_agent": "基于销售数据生成Word报告"
}
)
print("\n=== 最终输出 ===")
for worker_id, result in results.items():
print(f"{worker_id}: {result.status.value} -> {result.output[:50]}...")
# 运行演示
# asyncio.run(demo_master_worker())
主从模式的核心设计点:
- Master是循环的"调度器":它不负责具体业务逻辑,只负责"什么时候启动哪个Worker"
- Worker循环是独立的:每个Worker有自己的迭代、验证、超时机制
- Master等待策略:
asyncio.gather并行等待,但带整体超时控制
⚠️ 生产环境风险:Master单点故障会导致所有Worker不可控。建议生产环境补充Master冗余(主备切换)和Worker心跳检测机制。
验证步骤:
# 将上方完整代码(含demo_master_worker函数)保存为 master_worker.py,
# 在文件末尾添加一行 asyncio.run(demo_master_worker()),然后运行:
python master_worker.py
# 预期输出:
# [Master] 收到任务: 生成Q2销售报告
# [Master] 拆解为 2 个子任务
# [data_agent] 开始执行子任务: 查询Q2销售数据
# [doc_agent] 开始执行子任务: 基于销售数据生成Word报告
# [data_agent] ✅ 完成 (约500ms)
# [doc_agent] ✅ 完成 (约800ms)
# [Master] 任务汇总: 成功: 2/2
3.2 模式二:对等协调(Peer-to-Peer)
架构:多个Agent地位平等,通过协商达成一致,没有单一Master。
适用场景:
- 多Agent需要共同决策(如投票、共识)
- 没有单一权威,各Agent有独立信息源
- 需要容错(任一Agent故障不影响整体)
核心问题:多个Agent同时修改共享状态,怎么避免冲突?
Python实现(乐观锁协调):
import asyncio
from typing import Dict, Any, Optional
from dataclasses import dataclass, field
import time
@dataclass
class SharedState:
"""多Agent共享状态,带版本号(乐观锁)。"""
data: Dict[str, Any] = field(default_factory=dict)
version: int = 0
last_modified_by: str = ""
timestamp: float = field(default_factory=time.time)
class PeerToPeerCoordinator:
"""
对等协调器:多个Agent通过乐观锁协调共享状态。
GB/Z 185 8.1.3 冲突解决:乐观锁机制
"""
def __init__(self):
self.shared_state = SharedState()
self.state_lock = asyncio.Lock()
self.agent_commits: Dict[str, int] = {} # agent_id -> last_seen_version
async def read_state(self, agent_id: str) -> SharedState:
"""Agent读取共享状态(记录版本号用于后续校验)。"""
async with self.state_lock:
self.agent_commits[agent_id] = self.shared_state.version
return SharedState(
data=self.shared_state.data.copy(),
version=self.shared_state.version,
last_modified_by=self.shared_state.last_modified_by,
timestamp=self.shared_state.timestamp
)
async def commit_changes(self, agent_id: str,
expected_version: int,
changes: Dict[str, Any]) -> bool:
"""
Agent提交修改(乐观锁)。
Returns:
True: 提交成功
False: 版本冲突,需要重新读取后重试
"""
async with self.state_lock:
# 检查版本:如果共享状态已被其他Agent修改,拒绝提交
if self.shared_state.version != expected_version:
print(f" [{agent_id}] ❌ 版本冲突: 期望v{expected_version}, 实际v{self.shared_state.version}")
return False
# 应用修改
for key, value in changes.items():
self.shared_state.data[key] = value
self.shared_state.version += 1
self.shared_state.last_modified_by = agent_id
self.shared_state.timestamp = time.time()
self.agent_commits[agent_id] = self.shared_state.version
print(f" [{agent_id}] ✅ 提交成功: v{expected_version} -> v{self.shared_state.version}")
return True
async def propose_with_retry(self, agent_id: str,
proposal_generator: Callable,
max_retries: int = 3) -> bool:
"""
带重试的提案提交(解决冲突自动重试)。
proposal_generator: 函数,接收当前state,返回changes dict
"""
for attempt in range(max_retries):
# 1. 读取当前状态
state = await self.read_state(agent_id)
# 2. 基于当前状态生成提案
changes = proposal_generator(state.data)
# 3. 尝试提交
success = await self.commit_changes(agent_id, state.version, changes)
if success:
return True
print(f" [{agent_id}] 🔄 重试 {attempt + 1}/{max_retries}")
await asyncio.sleep(0.1 * (attempt + 1)) # 退避等待
print(f" [{agent_id}] ❌ 达到最大重试次数,提交失败")
return False
# ========== 使用示例:对等协调 ==========
async def demo_peer_to_peer():
"""演示:两个Agent同时修改共享状态。"""
coordinator = PeerToPeerCoordinator()
# 初始化共享状态
coordinator.shared_state.data = {"budget": 10000, "approved": False}
# Agent A提案:增加预算
def agent_a_proposal(current_data):
return {"budget": current_data.get("budget", 0) + 5000}
# Agent B提案:标记审批通过
def agent_b_proposal(current_data):
return {"approved": True}
# 两个Agent同时尝试提交
task_a = asyncio.create_task(
coordinator.propose_with_retry("agent_a", agent_a_proposal)
)
task_b = asyncio.create_task(
coordinator.propose_with_retry("agent_b", agent_b_proposal)
)
await asyncio.gather(task_a, task_b)
print(f"\n=== 最终共享状态 ===")
print(f"数据: {coordinator.shared_state.data}")
print(f"版本: v{coordinator.shared_state.version}")
print(f"最后修改者: {coordinator.shared_state.last_modified_by}")
# 运行演示
# asyncio.run(demo_peer_to_peer())
对等模式的核心设计点:
- 乐观锁:Agent先读取、再修改、最后提交时校验版本号
- 版本冲突自动重试:提交失败说明状态被其他Agent修改了,重新读取后重试
- 退避策略:重试间隔递增,避免多个Agent同时竞争
⚠️ 性能风险:乐观锁在冲突频繁时重试开销大,高竞争场景建议改用悲观锁(
asyncio.Lock)。Agent数量超过10个时,协商复杂度O(n²),建议改用主从模式。
验证步骤:
# 将上方完整代码(含demo_peer_to_peer函数)保存为 peer_coordinator.py,
# 在文件末尾添加一行 asyncio.run(demo_peer_to_peer()),然后运行:
python peer_coordinator.py
# 预期输出(Agent A和B的执行顺序可能交替):
# [agent_a] ✅ 提交成功: v0 -> v1
# [agent_b] ❌ 版本冲突: 期望v0, 实际v1
# [agent_b] 🔄 重试 1/3
# [agent_b] ✅ 提交成功: v1 -> v2
3.3 模式三:管道协调(Pipeline)
架构:Agent按顺序连接,前一个Agent的输出是下一个Agent的输入,像流水线一样。
适用场景:
- 任务有明确的先后顺序(如"先查数据→再生成报告→再校验")
- 每个Agent的输入严格依赖上一个Agent的输出
- 需要中间结果的质量控制(Verification Loop)
核心问题:前一个Agent的输出格式不符合下一个Agent的预期怎么办?
Python实现:
import inspect
from typing import Callable, List, Any, Optional
from dataclasses import dataclass
@dataclass
class PipelineStage:
"""管道的一个阶段。"""
stage_id: str
agent_func: Callable
input_schema: dict # 期望的输入格式
output_schema: dict # 输出的格式规范
verification_func: Optional[Callable] = None # 质量校验函数
class PipelineCoordinator:
"""
管道协调器:Agent按顺序执行,前一个的输出作为后一个的输入。
GB/Z 185 8.1.4 状态同步:通过管道中间状态传递
"""
def __init__(self):
self.stages: List[PipelineStage] = []
self.intermediate_results: Dict[str, Any] = {}
self.stage_outputs: Dict[str, Any] = {}
def add_stage(self, stage: PipelineStage):
"""添加一个管道阶段。"""
self.stages.append(stage)
async def execute(self, initial_input: Any) -> Any:
"""
按顺序执行管道中的所有阶段。
每个阶段:
1. 校验输入格式
2. 执行Agent逻辑
3. 校验输出质量(如果配置了Verification)
4. 如果质量不达标,触发重试或回滚
"""
current_input = initial_input
for stage in self.stages:
print(f"\n[Pipeline] 阶段 {stage.stage_id}: 开始")
# 1. 输入格式校验
if not self._validate_input(current_input, stage.input_schema):
error_msg = f"阶段 {stage.stage_id} 输入格式不符合预期"
print(f"[Pipeline] ❌ {error_msg}")
raise ValueError(error_msg)
# 2. 执行Agent逻辑
try:
if inspect.iscoroutinefunction(stage.agent_func):
output = await stage.agent_func(current_input)
else:
output = stage.agent_func(current_input)
except Exception as e:
print(f"[Pipeline] ❌ 阶段 {stage.stage_id} 执行失败: {e}")
raise
# 3. 输出质量校验(Verification Loop)
if stage.verification_func:
is_valid, feedback = stage.verification_func(output)
if not is_valid:
print(f"[Pipeline] ⚠️ 阶段 {stage.stage_id} 输出质量不达标: {feedback}")
# 策略:回退到上一个阶段或重试当前阶段
# 简化示例:直接抛出异常
raise ValueError(f"阶段 {stage.stage_id} 验证失败: {feedback}")
# 4. 保存中间结果
self.stage_outputs[stage.stage_id] = output
current_input = output # 作为下一个阶段的输入
print(f"[Pipeline] ✅ 阶段 {stage.stage_id} 完成")
return current_input
def _validate_input(self, input_data: Any, schema: dict) -> bool:
"""校验输入是否符合预期格式。"""
# 简化实现:实际应用jsonschema
if schema.get("type") == "string" and not isinstance(input_data, str):
return False
if schema.get("type") == "dict" and not isinstance(input_data, dict):
return False
return True
def get_pipeline_trace(self) -> List[Dict]:
"""获取管道执行的完整轨迹(用于审计)。"""
trace = []
for stage in self.stages:
trace.append({
"stage_id": stage.stage_id,
"input": self.intermediate_results.get(f"{stage.stage_id}_input"),
"output": self.stage_outputs.get(stage.stage_id)
})
return trace
# ========== 使用示例:管道协调 ==========
async def demo_pipeline():
"""演示:用户请求 -> 意图识别 -> 数据查询 -> 报告生成 -> 质量校验。"""
pipeline = PipelineCoordinator()
# 阶段1:意图识别Agent
def intent_agent(user_input: str) -> dict:
return {
"intent": "sales_report",
"quarter": "Q2",
"format": "word"
}
# 阶段2:数据查询Agent
def data_agent(intent_result: dict) -> dict:
return {
**intent_result,
"data": {"total": 12000000, "growth": "15%"},
"data_source": "sales_db"
}
# 阶段3:报告生成Agent
def report_agent(data_result: dict) -> str:
return f"Q2销售报告: 总销售额{data_result['data']['total']}元,环比增长{data_result['data']['growth']}"
# 阶段4:质量校验(Verification Loop)
def quality_check(report: str) -> tuple:
has_numbers = any(c.isdigit() for c in report)
has_growth = "增长" in report or "环比" in report
return has_numbers and has_growth, "缺少关键数据" if not has_numbers else "OK"
# 构建管道
pipeline.add_stage(PipelineStage(
stage_id="intent",
agent_func=intent_agent,
input_schema={"type": "string"},
output_schema={"type": "dict"}
))
pipeline.add_stage(PipelineStage(
stage_id="data_query",
agent_func=data_agent,
input_schema={"type": "dict"},
output_schema={"type": "dict"}
))
pipeline.add_stage(PipelineStage(
stage_id="report",
agent_func=report_agent,
input_schema={"type": "dict"},
output_schema={"type": "string"},
verification_func=quality_check
))
# 执行
result = await pipeline.execute("帮我生成Q2销售报告")
print(f"\n=== 最终输出 ===\n{result}")
# 审计轨迹
trace = pipeline.get_pipeline_trace()
print(f"\n=== 执行轨迹 ===")
for step in trace:
print(f" {step['stage_id']}: {step['output']}")
# 运行演示
# asyncio.run(demo_pipeline())
管道模式的核心设计点:
- Schema约束:每个阶段的输入/输出都有格式要求,防止数据在管道中"变形"
- Verification Loop嵌入:每个阶段完成后可以自动质量校验,不达标就阻断
- 审计轨迹完整:记录每个阶段的输入/输出,问题定位时可以追溯
⚠️ 故障传递风险:中间任一阶段失败会导致整个管道中断。建议:①每个阶段设置独立超时;②关键阶段后加Checkpoint便于恢复;③失败时支持跳过阶段(Skippable Stage)或使用降级数据。
验证步骤:
# 将上方完整代码(含demo_pipeline函数)保存为 pipeline_demo.py,
# 在文件末尾添加一行 asyncio.run(demo_pipeline()),然后运行:
python pipeline_demo.py
# 输入:'帮我生成Q2销售报告'
# 预期输出:
# [Pipeline] 阶段 intent: 开始
# [Pipeline] ✅ 阶段 intent 完成
# [Pipeline] 阶段 data_query: 开始
# [Pipeline] ✅ 阶段 data_query 完成
# [Pipeline] 阶段 report: 开始
# [Pipeline] ✅ 阶段 report 完成
# === 最终输出 ===
# Q2销售报告: 总销售额12000000元,环比增长15%
四、四大核心问题的统一解决方案
4.1 同步问题:Agent A怎么知道Agent B完成了?
三种模式的解决方案对比:
| 模式 | 同步机制 | 代码实现 |
|---|---|---|
| 主从 | Master用asyncio.gather等待所有Worker | await asyncio.gather(*tasks) |
| 对等 | 版本号通知(读取时获取最新版本) | shared_state.version |
| 管道 | 前一阶段完成自动触发下一阶段 | current_input = output |
统一建议:使用事件总线(Event Bus)作为同步基础设施:
import inspect
from typing import Dict, List, Callable, Any
class AgentEventBus:
"""Agent间事件总线,用于循环同步。"""
def __init__(self):
self.subscribers: Dict[str, List[Callable]] = {}
self.event_queue = 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 publish(self, event_type: str, data: Any):
"""发布事件。"""
await self.event_queue.put({"type": event_type, "data": data})
# 通知订阅者
for callback in self.subscribers.get(event_type, []):
if inspect.iscoroutinefunction(callback):
asyncio.create_task(callback(data))
else:
callback(data)
# 事件总线使用示例
async def demo_event_bus():
bus = AgentEventBus()
async def on_task_done(data):
print(f"[订阅者] 收到任务完成事件: {data}")
bus.subscribe("task.completed", on_task_done)
await bus.publish("task.completed", {"task_id": "t1", "result": "ok"})
await asyncio.sleep(0.1)
# asyncio.run(demo_event_bus())
# 预期输出:
# [订阅者] 收到任务完成事件: {'task_id': 't1', 'result': 'ok'}
4.2 冲突问题:两个Agent同时修改数据怎么办?
| 模式 | 冲突解决策略 | 适用场景 |
|---|---|---|
| 主从 | 不存在冲突(Master单点控制) | 任务可明确分配 |
| 对等 | 乐观锁+自动重试 | 需要并行协商 |
| 管道 | 不存在冲突(串行执行) | 任务有明确顺序 |
4.3 超时问题:Agent B卡住了怎么办?
统一策略:
# 在MasterWorkerCoordinator中已实现
await asyncio.wait_for(
asyncio.gather(*tasks, return_exceptions=True),
timeout=self.default_timeout
)
超时后的处理策略:
- 标记为TIMEOUT:记录状态,不阻塞其他Agent
- 降级处理:用缓存数据或默认值替代
- 重试机制:指数退避后重试(最多3次)
- 告警通知:发送超时告警给运维系统
4.4 状态一致性问题:Agent A重启了,怎么恢复?
解决方案:Checkpoint + 事件溯源(Event Sourcing)
import os, json, time
from typing import Optional
class AgentStateManager:
"""Agent状态管理:支持Checkpoint和恢复。"""
def __init__(self, storage_path: str = "agent_states"):
self.storage_path = storage_path
os.makedirs(storage_path, exist_ok=True)
def save_checkpoint(self, agent_id: str, state: dict):
"""保存Agent状态快照。"""
filepath = os.path.join(self.storage_path, f"{agent_id}_checkpoint.json")
with open(filepath, "w") as f:
json.dump({"agent_id": agent_id, "state": state, "timestamp": time.time()}, f)
def load_checkpoint(self, agent_id: str) -> Optional[dict]:
"""加载Agent状态快照。"""
filepath = os.path.join(self.storage_path, f"{agent_id}_checkpoint.json")
if os.path.exists(filepath):
with open(filepath, "r") as f:
return json.load(f)
return None
# ========== 使用示例 ==========
# manager = AgentStateManager()
# manager.save_checkpoint("agent_a", {"iteration": 5, "results": ["r1", "r2"]})
# checkpoint = manager.load_checkpoint("agent_a")
# print(checkpoint["state"]["iteration"]) # 预期输出:5
五、与GB/Z 185和A2A的融合:标准→协议→代码
本文是标准、协议、代码三层的融合实现:
GB/Z 185 8.1(标准层)
├── 8.1.1 协作协议 → A2A协议实现
├── 8.1.2 角色分工 → 主从/对等/管道三种模式
├── 8.1.3 冲突解决 → 乐观锁/仲裁者
└── 8.1.4 状态同步 → Checkpoint+事件溯源
↓
A2A协议(协议层)
├── Agent Card → Worker能力注册
├── Send Task → Master分配任务给Worker
├── Get Task → Master查询Worker状态
└── Cancel Task → Master超时后取消Worker
↓
Loop Engineering(代码层)
├── MasterWorkerCoordinator → 主从循环协调
├── PeerToPeerCoordinator → 对等循环协调
└── PipelineCoordinator → 管道循环协调
六、方案对比与选型建议
6.1 三种模式综合对比
| 维度 | 主从(Master-Worker) | 对等(Peer-to-Peer) | 管道(Pipeline) |
|---|---|---|---|
| 拓扑结构 | 星型(Master中心控制) | 网状(全互联) | 链型(线性串联) |
| 控制流 | 集中式(Master统一调度) | 去中心化(协商共识) | 顺序式(前驱驱动) |
| 数据流 | 分散→汇总 | 共享状态+版本控制 | 前输出→后输入 |
| 冲突风险 | ✅ 无冲突(单点分配) | ⚠️ 需乐观锁/悲观锁 | ✅ 无冲突(串行) |
| 容错能力 | ⚠️ Master单点故障 | ✅ 任一节点故障不影响整体 | ⚠️ 中间阶段故障整体中断 |
| 扩展性 | ⭐⭐⭐ 增加Worker简单 | ⭐⭐ 协商复杂度随节点数上升 | ⭐ 新增阶段需对齐Schema |
| 典型场景 | 任务分发+结果汇总 | 多Agent投票/共识决策 | 数据处理流水线 |
6.2 选型决策指南
6.3 适用边界说明
⚠️ 主从模式不适用场景:Worker之间需要直接通信、Master成为性能瓶颈、高可用要求严格(Master单点故障)
⚠️ 对等模式不适用场景:Agent数量>10个(协商复杂度O(n²))、需要严格的事务一致性、实时性要求极高
⚠️ 管道模式不适用场景:阶段之间有循环依赖、某个阶段处理时间不可预测长、需要动态增减阶段
6.4 快速选型卡(30秒判断用哪种)
主从模式 ← 任务可拆解、需要统一汇总
[ ] 子任务可以并行执行
[ ] 需要Master做最终决策
[ ] Worker之间不需要直接通信
[ ] 能接受Master单点故障风险
→ 不符合任一条件 → 看对等或管道模式
对等模式 ← 需要共识、没有单一权威
[ ] Agent数量 ≤ 10(超过10个协商复杂度O(n²))
[ ] 需要容错(任一Agent故障不影响整体)
[ ] 不要求严格的事务一致性
[ ] 能容忍偶尔的版本冲突重试
→ 冲突率超过30% → 建议换悲观锁或主从模式
管道模式 ← 任务有明确顺序、需要质量控制
[ ] 前一阶段输出是后一阶段输入(线性依赖)
[ ] 需要中间结果校验(Verification Loop)
[ ] 要求完整审计轨迹
[ ] 能接受中间阶段故障导致整体中断
→ 阶段之间有循环依赖 → 不适合管道,用对等+事件驱动
6.5 混合使用策略
实际生产系统中,三种模式往往组合使用:
# 混合架构示例:Master-Worker + Pipeline
# 调度Agent(Master)按管道顺序调度Worker
# 每个Worker内部又有自己的Pipeline
async def hybrid_demo():
coordinator = MasterWorkerCoordinator(default_timeout=30.0)
# Worker1: 数据采集管道(内部3个串行阶段)
async def data_worker(task: str):
pipeline = PipelineCoordinator()
pipeline.add_stage(PipelineStage("extract", ...))
pipeline.add_stage(PipelineStage("transform", ...))
pipeline.add_stage(PipelineStage("load", ...))
return await pipeline.execute(task)
# Worker2: 分析Agent(对等协商模式)
async def analysis_worker(task: str):
peer = PeerToPeerCoordinator()
return await peer.propose_with_retry("analysis_agent", lambda d: d)
coordinator.register_worker("data", data_worker)
coordinator.register_worker("analysis", analysis_worker)
return await coordinator.execute_task("hybrid", ...)
七、总结
本文从GB/Z 185标准出发,结合A2A协议和Loop Engineering,实现了三种Multi-Agent循环协调模式:
| 模式 | 核心问题 | 解决方案 | 选型信号 |
|---|---|---|---|
| 主从 | Worker循环怎么被Master调度 | asyncio.gather + 超时控制 | 任务可拆解、需统一汇总 |
| 对等 | 共享状态冲突怎么解决 | 乐观锁 + 版本号 + 自动重试 | 需要共识、无单一权威 |
| 管道 | 前一阶段输出不符合下一阶段输入 | Schema校验 + Verification Loop | 任务有明确顺序、需质量控制 |
核心结论:
- Multi-Agent不是多个单Agent的简单叠加,而是分布式系统,需要专门的循环协调机制
- GB/Z 185 8.1条款是设计Multi-Agent系统的"检查清单"——角色分工、冲突解决、状态同步,每个都有对应工程实现
- A2A协议是循环协调的通信层——Agent Card注册能力、Task传递任务、Cancel处理超时
- 三种模式可以混合使用——实际系统中,Master-Worker内部可能包含Pipeline,Pipeline某个阶段可能用对等协商
相关阅读:
- GB/Z 185合规的Agent Loop设计(标准合规基础——本文的上篇)
- A2A协议实战:用Python实现Agent-to-Agent通信(A2A协议——Multi-Agent的通信层)
- Loop Engineering四层架构(Loop控制基础——单Agent循环的工程化)
- MCP Server安全加固(安全层——Multi-Agent间的权限控制)
你现在的多Agent系统用的哪种协调方式?
- A. 主从模式(一个Master统一调度)
- B. 消息队列(RabbitMQ/Kafka做中间人)
- C. 自研协调器
- D. 还没用到多Agent,先收藏
——评论区说说你碰到最多的坑:同步阻塞、数据冲突、超时无响应、还是状态丢了拉不起来?哪类问题呼声最高,我下篇出对应的排错指南+完整代码。
觉得有用的话点赞+收藏,设计多智能体系统时直接翻出这三种模式的代码来参考,让更多正在做Multi-Agent的开发者看到这四种解决思路。
📅 更新日志:
2026-07-21 · 1.0版发布
⚠️ 版本变更提示:本文基于Python 3.10+ asyncio标准库实现。如果Python大版本升级(如3.13+),请关注asyncio官方迁移指南。三种协调模式的算法逻辑不受Python版本影响,可直接移植到其他语言(Go/Java/TypeScript)。
更多推荐


所有评论(0)