多 Agent 的指挥系统:Agent 编排引擎的设计与实现

一、单 Agent 的编排瓶颈:为什么"一个 Agent 干所有事"走不通

单个 LLM Agent 处理简单任务没问题,但面对多步骤、多工具的复杂工作流时,问题就来了:Prompt 越来越长,工具描述互相干扰,Agent 在不同角色间反复跳转,输出质量断崖式下降。更致命的是,单 Agent 无法并行——一个需要同时查询数据库和调用外部 API 的任务,只能串行执行,延迟翻倍。

Agent 编排引擎解决的就是这个问题:把复杂工作流拆解为多个 Agent 的协作任务,由编排引擎负责任务分发、依赖管理、并行调度和结果聚合。每个 Agent 只关心自己的领域,上下文干净,工具集精简,输出质量稳定。编排引擎是"指挥",Agent 是"乐手"——指挥不需要会弹每件乐器,但必须知道什么时候谁该进、谁该停。

二、Agent 编排引擎的核心架构

flowchart TD
    A[工作流定义: DAG] --> B[编排引擎: 解析 + 调度]
    B --> C[Agent A: 数据查询]
    B --> D[Agent B: API 调用]
    B --> E[Agent C: 文档生成]

    C -->|结果| F[结果聚合器]
    D -->|结果| F
    F --> G[Agent D: 质量检查]
    G -->|通过| H[最终输出]
    G -->|不通过| I[重试/降级]

    subgraph 编排核心
        J[依赖解析: 拓扑排序]
        K[并行调度: 无依赖任务并发执行]
        L[超时控制: 单 Agent 超时熔断]
        M[重试策略: 指数退避 + 最大次数]
    end

    subgraph Agent 隔离
        N[独立上下文: 各 Agent 只看到自己的工具和知识]
        O[独立资源: Token 预算独立分配]
        P[独立超时: 不因一个 Agent 卡住整体]
    end

编排引擎的三个核心能力

依赖解析:工作流定义为有向无环图(DAG),引擎通过拓扑排序确定执行顺序。无依赖关系的节点并行执行,有依赖关系的节点串行执行。

并行调度:并行度受限于 LLM API 的并发限制和 Token 预算。引擎需要实现信号量机制,控制同时运行的 Agent 数量。

容错与降级:单个 Agent 失败不应导致整个工作流崩溃。引擎需要实现超时熔断、重试和降级策略——如果某个 Agent 超时,用缓存结果或默认值替代,而非阻塞整个流程。

三、Agent 编排引擎的实现

3.1 工作流定义与 DAG 解析

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


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


@dataclass
class AgentTask:
    """编排引擎中的单个 Agent 任务"""
    id: str
    agent_name: str
    prompt_template: str
    dependencies: list[str] = field(default_factory=list)
    timeout_seconds: int = 60
    max_retries: int = 2
    fallback_value: Any = None
    state: AgentState = AgentState.PENDING
    result: Any = None
    error: Optional[str] = None


@dataclass
class Workflow:
    """工作流定义:DAG 结构"""
    name: str
    tasks: dict[str, AgentTask]
    max_parallelism: int = 3


class DAGResolver:
    """DAG 依赖解析与拓扑排序"""

    def topological_sort(self, workflow: Workflow) -> list[list[str]]:
        """返回分层执行计划:同一层内的任务可并行执行"""
        in_degree = {tid: 0 for tid in workflow.tasks}
        for task in workflow.tasks.values():
            for dep in task.dependencies:
                in_degree[task.id] += 1

        layers = []
        remaining = set(workflow.tasks.keys())

        while remaining:
            # 找到当前无依赖的任务
            ready = [tid for tid in remaining if in_degree[tid] == 0]
            if not ready:
                raise ValueError("工作流存在循环依赖")

            layers.append(ready)
            for tid in ready:
                remaining.remove(tid)
                # 减少依赖此任务的后续任务的入度
                for other_id in remaining:
                    if tid in workflow.tasks[other_id].dependencies:
                        in_degree[other_id] -= 1

        return layers

3.2 编排引擎核心调度

import asyncio
from datetime import datetime


class OrchestrationEngine:
    """Agent 编排引擎"""

    def __init__(self, agent_registry: dict, max_parallelism: int = 3):
        self.agent_registry = agent_registry  # agent_name → agent_instance
        self.max_parallelism = max_parallelism
        self.semaphore = asyncio.Semaphore(max_parallelism)

    async def execute(self, workflow: Workflow, context: dict) -> dict:
        """执行完整工作流"""
        resolver = DAGResolver()
        layers = resolver.topological_sort(workflow)

        results = {}
        start_time = datetime.now()

        for layer_idx, layer in enumerate(layers):
            # 同一层内的任务并行执行
            tasks = []
            for task_id in layer:
                task = workflow.tasks[task_id]
                # 构建任务输入:依赖任务的结果
                task_input = self._build_task_input(task, results, context)
                tasks.append(self._execute_task(task, task_input))

            layer_results = await asyncio.gather(*tasks, return_exceptions=True)

            # 收集结果
            for task_id, result in zip(layer, layer_results):
                if isinstance(result, Exception):
                    workflow.tasks[task_id].state = AgentState.FAILED
                    workflow.tasks[task_id].error = str(result)
                    results[task_id] = workflow.tasks[task_id].fallback_value
                else:
                    results[task_id] = result

        total_time = (datetime.now() - start_time).total_seconds()

        return {
            "workflow": workflow.name,
            "results": results,
            "total_time_seconds": total_time,
            "task_states": {
                tid: t.state.value for tid, t in workflow.tasks.items()
            },
        }

    async def _execute_task(self, task: AgentTask, task_input: dict) -> Any:
        """执行单个 Agent 任务,带超时和重试"""
        async with self.semaphore:
            agent = self.agent_registry.get(task.agent_name)
            if not agent:
                raise ValueError(f"未注册的 Agent: {task.agent_name}")

            for attempt in range(task.max_retries + 1):
                try:
                    task.state = AgentState.RUNNING
                    result = await asyncio.wait_for(
                        agent.run(task_input),
                        timeout=task.timeout_seconds,
                    )
                    task.state = AgentState.SUCCESS
                    task.result = result
                    return result

                except asyncio.TimeoutError:
                    task.state = AgentState.TIMEOUT
                    if attempt < task.max_retries:
                        await asyncio.sleep(2 ** attempt)  # 指数退避
                        continue
                    return task.fallback_value

                except Exception as e:
                    task.state = AgentState.FAILED
                    task.error = str(e)
                    if attempt < task.max_retries:
                        await asyncio.sleep(2 ** attempt)
                        continue
                    return task.fallback_value

    def _build_task_input(
        self, task: AgentTask, results: dict, context: dict
    ) -> dict:
        """构建任务输入:合并上下文和依赖结果"""
        task_input = {"context": context}
        for dep_id in task.dependencies:
            if dep_id in results and results[dep_id] is not None:
                task_input[f"dep_{dep_id}"] = results[dep_id]
        return task_input

3.3 工作流定义示例

# 定义一个"竞品分析报告生成"工作流
competitor_analysis_workflow = Workflow(
    name="competitor_analysis",
    max_parallelism=3,
    tasks={
        "fetch_data": AgentTask(
            id="fetch_data",
            agent_name="web_search_agent",
            prompt_template="搜索 {company_name} 的最新产品动态和用户评价",
            timeout_seconds=30,
            max_retries=2,
            fallback_value={"search_results": "数据获取失败,使用缓存数据"},
        ),
        "analyze_sentiment": AgentTask(
            id="analyze_sentiment",
            agent_name="sentiment_agent",
            prompt_template="分析以下用户评价的情感倾向",
            dependencies=["fetch_data"],
            timeout_seconds=45,
        ),
        "fetch_financials": AgentTask(
            id="fetch_financials",
            agent_name="api_agent",
            prompt_template="获取 {company_name} 的财务数据",
            timeout_seconds=20,
        ),
        "generate_report": AgentTask(
            id="generate_report",
            agent_name="writing_agent",
            prompt_template="基于情感分析和财务数据生成竞品分析报告",
            dependencies=["analyze_sentiment", "fetch_financials"],
            timeout_seconds=60,
        ),
        "quality_check": AgentTask(
            id="quality_check",
            agent_name="review_agent",
            prompt_template="检查报告的完整性和准确性",
            dependencies=["generate_report"],
            timeout_seconds=30,
        ),
    },
)

四、编排引擎的工程权衡

并行度 vs Token 预算:并行执行多个 Agent 时,每个 Agent 的上下文独立,总 Token 消耗是各 Agent 之和。如果工作流有 5 个并行 Agent,每个消耗 2000 Token,总消耗就是 10000 Token。必须设置全局 Token 预算上限,超限时降级为串行执行。

降级策略的精度损失:使用 fallback_value 替代失败 Agent 的结果,意味着最终输出可能基于不完整数据。编排引擎必须在输出中标注哪些数据是降级结果,让下游消费者知道数据的可靠性。

工作流定义的灵活性:静态 DAG 无法处理条件分支(如"如果情感分析结果为负面,增加危机分析步骤")。更高级的编排引擎需要支持动态 DAG——根据中间结果调整后续任务的执行计划。但这增加了引擎的复杂度,也增加了调试难度。

可观测性:多 Agent 工作流的调试比单 Agent 困难得多。必须为每个任务记录输入、输出、耗时和状态,建立类似分布式链路追踪的 Agent Trace 机制。

五、总结

Agent 编排引擎是多 Agent 系统的"指挥系统",核心能力是依赖解析、并行调度和容错降级。DAG 定义工作流结构,拓扑排序确定执行层次,信号量控制并行度,超时和重试保障可靠性。

落地要点:工作流用 DAG 定义,无依赖任务并行执行;每个 Agent 独立上下文和超时,互不干扰;失败任务用 fallback 降级,不阻塞整体流程;全链路 Trace 是调试的基础。编排引擎的目标不是让 Agent 做更多事,而是让多个 Agent 做各自最擅长的事——专业分工,高效协作。

Logo

小龙虾开发者社区是 CSDN 旗下专注 OpenClaw 生态的官方阵地,聚焦技能开发、插件实践与部署教程,为开发者提供可直接落地的方案、工具与交流平台,助力高效构建与落地 AI 应用

更多推荐