多 Agent 协作编排:从任务分解到执行监控的工程实践

cover

一、单 Agent 的能力天花板:复杂任务的协作困境

当 AI Agent 从概念验证走向生产环境时,一个核心瓶颈迅速暴露:单 Agent 处理复杂任务时,上下文窗口溢出、工具调用链过长、错误恢复困难。以一个典型的"技术调研报告生成"任务为例,单 Agent 需要同时负责信息检索、内容筛选、结构编排、格式校验和引用核查——每个子任务的认知负载叠加后,输出质量急剧下降。

多 Agent 协作的核心思路是将复杂任务拆解为职责单一的子任务,每个 Agent 只负责一个明确的工作域。但多 Agent 系统引入了新的工程挑战:Agent 间的通信协议如何设计?任务依赖关系如何编排?部分 Agent 执行失败时如何降级?这些问题决定了系统从 Demo 到生产的距离。

二、多 Agent 编排的核心架构与通信机制

多 Agent 系统的架构设计需要解决三个核心问题:任务分解、通信协议和执行监控。

flowchart TB
    subgraph 任务分解层
        A[用户请求] --> B[Planner Agent]
        B --> C[任务 DAG 生成]
        C --> D[子任务 1: 信息检索]
        C --> E[子任务 2: 数据分析]
        C --> F[子任务 3: 内容生成]
    end

    subgraph 通信与执行层
        D --> G[消息总线]
        E --> G
        F --> G
        G --> H[共享状态存储]
        H --> I[结果聚合器]
    end

    subgraph 监控与容错层
        G --> J[执行追踪器]
        J --> K{超时/失败?}
        K -->|是| L[降级策略]
        K -->|否| M[继续执行]
        L --> N[重试/跳过/人工介入]
        I --> O[最终输出]
    end

任务分解由 Planner Agent 完成,它将用户请求解析为有向无环图(DAG),每个节点代表一个子任务,边代表依赖关系。关键设计决策:DAG 的粒度不宜过细——每个子任务应是一个 Agent 可以独立完成的完整工作单元,而非单步操作。

通信协议采用消息总线模式,Agent 之间不直接调用,而是通过消息总线传递结构化消息。这解耦了 Agent 的执行时序,同时支持异步并行。消息格式必须包含:发送者 ID、接收者 ID、消息类型(请求/响应/事件)、负载内容和关联的任务 ID。

执行监控通过执行追踪器实现,记录每个子任务的开始时间、结束时间、状态和输出摘要。当检测到超时或失败时,触发降级策略:重试(可恢复的临时错误)、跳过(非关键路径的子任务)或人工介入(关键路径的不可恢复错误)。

三、生产级多 Agent 编排框架的实现

import asyncio
import uuid
from dataclasses import dataclass, field
from enum import Enum
from typing import Any, Callable, Optional


class TaskStatus(Enum):
    PENDING = "pending"
    RUNNING = "running"
    SUCCESS = "success"
    FAILED = "failed"
    SKIPPED = "skipped"


@dataclass
class AgentMessage:
    """Agent 间通信的结构化消息"""
    sender: str
    receiver: str
    msg_type: str  # request / response / event
    payload: dict
    task_id: str
    correlation_id: str = field(default_factory=lambda: uuid.uuid4().hex[:8])


@dataclass
class SubTask:
    """子任务定义"""
    task_id: str
    agent_name: str
    handler: Callable
    dependencies: list[str] = field(default_factory=list)
    timeout: float = 60.0
    max_retries: int = 1
    # 非关键路径的子任务失败后可跳过
    critical: bool = True


class MessageBus:
    """异步消息总线,解耦 Agent 间的直接依赖"""

    def __init__(self):
        self._subscribers: dict[str, list[asyncio.Queue]] = {}

    def subscribe(self, agent_name: str) -> asyncio.Queue:
        if agent_name not in self._subscribers:
            self._subscribers[agent_name] = []
        queue: asyncio.Queue = asyncio.Queue()
        self._subscribers[agent_name].append(queue)
        return queue

    async def publish(self, message: AgentMessage):
        subscribers = self._subscribers.get(message.receiver, [])
        for queue in subscribers:
            await queue.put(message)


class MultiAgentOrchestrator:
    """多 Agent 编排引擎"""

    def __init__(self):
        self.bus = MessageBus()
        self.shared_state: dict[str, Any] = {}
        self.task_status: dict[str, TaskStatus] = {}
        self._execution_log: list[dict] = []

    def _log(self, task_id: str, event: str, detail: str = ""):
        self._execution_log.append({
            "task_id": task_id,
            "event": event,
            "detail": detail,
        })

    async def _execute_task(self, task: SubTask) -> Any:
        """执行单个子任务,含超时与重试"""
        self.task_status[task.task_id] = TaskStatus.RUNNING
        self._log(task.task_id, "start")

        for attempt in range(task.max_retries + 1):
            try:
                result = await asyncio.wait_for(
                    task.handler(self.shared_state, self.bus),
                    timeout=task.timeout,
                )
                self.task_status[task.task_id] = TaskStatus.SUCCESS
                self.shared_state[task.task_id] = result
                self._log(task.task_id, "success")
                return result
            except asyncio.TimeoutError:
                self._log(task.task_id, "timeout", f"attempt={attempt}")
            except Exception as e:
                self._log(task.task_id, "error", f"attempt={attempt}, err={e}")

            if attempt < task.max_retries:
                await asyncio.sleep(1.0 * (2 ** attempt))

        # 重试耗尽
        if task.critical:
            self.task_status[task.task_id] = TaskStatus.FAILED
            raise RuntimeError(f"关键任务 {task.task_id} 执行失败")
        else:
            self.task_status[task.task_id] = TaskStatus.SKIPPED
            self._log(task.task_id, "skipped", "non-critical task failed")
            return None

    async def run(self, tasks: list[SubTask]) -> dict:
        """按 DAG 依赖关系编排执行"""
        # 初始化所有任务状态
        for t in tasks:
            self.task_status[t.task_id] = TaskStatus.PENDING

        completed: set[str] = set()
        max_iterations = len(tasks) * 2  # 防止死循环

        for _ in range(max_iterations):
            # 找出所有依赖已满足的待执行任务
            ready = [
                t for t in tasks
                if self.task_status[t.task_id] == TaskStatus.PENDING
                and all(d in completed for d in t.dependencies)
            ]

            if not ready:
                if len(completed) == len(tasks):
                    break
                # 存在无法满足依赖的任务,标记为失败
                for t in tasks:
                    if self.task_status[t.task_id] == TaskStatus.PENDING:
                        self.task_status[t.task_id] = TaskStatus.FAILED
                break

            # 并行执行所有就绪任务
            coros = [self._execute_task(t) for t in ready]
            results = await asyncio.gather(*coros, return_exceptions=True)

            for t, r in zip(ready, results):
                if isinstance(r, Exception):
                    if t.critical:
                        # 关键任务失败,终止整个编排
                        return self._build_result()
                completed.add(t.task_id)

        return self._build_result()

    def _build_result(self) -> dict:
        return {
            "status": self.task_status,
            "shared_state": self.shared_state,
            "execution_log": self._execution_log,
        }

关键设计点:

  1. 消息总线解耦:Agent 之间通过 MessageBus 通信,而非直接函数调用。这使得 Agent 可以独立部署和扩展,也为后续引入远程 Agent(如通过 HTTP 调用的第三方服务)预留了接口。

  2. DAG 依赖调度:编排引擎按拓扑顺序执行任务,同一层级的任务并行执行。依赖检测通过 completed 集合实现,避免轮询开销。

  3. 关键路径保护critical 标志区分关键和非关键任务。非关键任务失败后跳过,不影响整体流程;关键任务失败则终止编排。

四、多 Agent 系统的工程代价与适用边界

通信开销不可忽视。Agent 间的消息传递引入了序列化/反序列化成本和异步等待延迟。在子任务执行时间小于 100ms 的场景下,通信开销可能占总耗时的 50% 以上。此时单 Agent 方案反而更高效。

状态一致性是隐性地雷。共享状态存储在并发写入时可能出现竞态条件。上述实现使用字典存储,在单进程内是安全的,但一旦扩展到多进程或多节点,就需要引入分布式锁或 CRDT(无冲突复制数据类型)来解决一致性问题。

调试难度指数级增长。当 5 个 Agent 协作完成一个任务时,一次失败可能涉及 3 个 Agent 的交互链路。执行日志(_execution_log)是排障的唯一线索,必须包含完整的消息流转记录。

适用边界:多 Agent 编排适合"子任务之间有明确边界、需要不同工具或模型、执行时间在秒级以上"的场景。对于"子任务之间高度耦合、执行时间在毫秒级"的场景,单 Agent 方案更优。

五、总结

多 Agent 协作编排的核心价值在于将复杂任务的认知负载分散到职责单一的 Agent 上,通过 DAG 依赖调度实现并行执行和容错降级。落地时需把握三个要点:第一,任务分解粒度以"一个 Agent 可独立完成的工作单元"为标准,过细的粒度反而增加通信开销;第二,Agent 间通过消息总线通信而非直接调用,为后续扩展和独立部署预留空间;第三,区分关键路径和非关键路径,非关键任务失败后跳过而非终止,提升系统的韧性。多 Agent 不是万能架构,在子任务高度耦合或执行时间极短的场景下,单 Agent 方案仍然是更务实的选择。

Logo

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

更多推荐