多智能体互联时的循环协调: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重启了,怎么恢复到之前的状态?(状态一致性问题)
文档Agent (Worker B) 数据Agent (Worker A) 调度Agent (Coordinator) 文档Agent (Worker B) 数据Agent (Worker A) 调度Agent (Coordinator) 卡在第5轮迭代 没有响应 缺少数据输入 空转等待 发送任务(A2A) Loop执行中... 发送任务(A2A) 开始执行... 超时?重试?还是放弃?

核心洞察: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)的循环节奏。

返回结果

数据Agent
Worker Loop 1

文档Agent
Worker Loop 2

邮件Agent
Worker Loop 3

适用场景

  • 任务可以被明确拆解为子任务
  • 需要统一调度、统一汇总
  • 各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 A

Agent B

Agent C

适用场景

  • 多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 1
意图识别

Agent 2
数据查询

Agent 3
结果生成

Agent 4
质量校验

最终输出

适用场景

  • 任务有明确的先后顺序(如"先查数据→再生成报告→再校验")
  • 每个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等待所有Workerawait 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
)

超时后的处理策略

  1. 标记为TIMEOUT:记录状态,不阻塞其他Agent
  2. 降级处理:用缓存数据或默认值替代
  3. 重试机制:指数退避后重试(最多3次)
  4. 告警通知:发送超时告警给运维系统

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 选型决策指南

无依赖

Agent间有依赖关系?

依赖是线性的?

管道模式

需要共享状态?

有统一调度需求?

主从模式

对等模式

对等模式+乐观锁

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任务有明确顺序、需质量控制

核心结论

  1. Multi-Agent不是多个单Agent的简单叠加,而是分布式系统,需要专门的循环协调机制
  2. GB/Z 185 8.1条款是设计Multi-Agent系统的"检查清单"——角色分工、冲突解决、状态同步,每个都有对应工程实现
  3. A2A协议是循环协调的通信层——Agent Card注册能力、Task传递任务、Cancel处理超时
  4. 三种模式可以混合使用——实际系统中,Master-Worker内部可能包含Pipeline,Pipeline某个阶段可能用对等协商

相关阅读:


你现在的多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)。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐