摘要:2026 年 AI Agent 已从"单兵作战"进入"智能体集群"时代。但生产落地最大的坑不在框架选择,而在三个工程基座的质量:任务拆解是否产生可追踪的子任务 DAG、工具调用是否容忍部分失败、记忆系统是否跨会话一致。本文从编排模式选型出发,系统展开任务拆解(层级规划 + ReAct 动态补全)、工具调用(三级容错闭环)、记忆模块(Working/Episodic/Semantic 三层架构)三大核心设计,附 LangGraph Supervisor 完整可运行代码与一线避坑清单。读完你会抛弃"单 Agent 装所有工具"的反模式,建立起可演进的 Multi-Agent 生产架构。


目录


一、四种编排模式:什么时候该用 Multi-Agent?

1.1 单 Agent 的天花板

把 20 个工具挂在一个 Agent 上,用一条超长 System Prompt 描述所有场景——这是 2025 年最常见的 Agent 反模式。微软 Azure 架构中心 2026 年发布的 Agent 编排指南明确指出:单 Agent 的可靠性会随着工具数量增加而显著下降——每个新增工具都在扩大 LLM 在每个推理步骤中需要同时考虑的范围(工具的 Schema、当前上下文、调用历史),导致工具选择准确率持续退化。这正是将复杂系统拆分为 Multi-Agent 的根本动机。

何时单 Agent 够用、何时必须 Multi-Agent?微软 Azure 在 2026 年发布的 Agent 编排指南中给出了决策框架:

复杂度等级 方案 适用场景
L1 单次 LLM 调用 分类、摘要、翻译等单步任务
L2 单 Agent + 工具 同一领域内的动态工具选择(如订单查询、数据库检索)
L3 Multi-Agent 编排 跨领域/跨职能、需独立安全边界、需并行专业化

决策原则:用能满足需求的最低复杂度。 如果你只是做一个 RAG 问答,不需要 Multi-Agent。当你的系统涉及"研发 + 测试 + 审批"三个不同职能、且各有不同工具和权限时,才上 Multi-Agent。

1.2 四种编排模式对比

2026 年生产环境验证的四种 Multi-Agent 编排模式:

模式 核心机制 代表框架 适用场景 关键风险
Sequential Pipeline 固定阶段,数据单向流动 CrewAI Sequential 已知步骤的流水线(分析→规划→执行→验证) 不支持动态分支,中间失败需人工重启
Supervisor/Orchestrator 一个管理者 Agent 分配任务给专业 Worker LangGraph Supervisor 任务类型多变、需动态路由 Supervisor 节点是单点,需加超时+降级
Swarm/Handoff Agent 间平等协商,按上下文自主转交 OpenAI Swarm 对话分流、客服路由 循环调用难以检测,调试困难
Adversarial/Review 生成 Agent + 批判 Agent,多轮对抗提升质量 AutoGen 代码审查、报告生成、安全审计 轮次不可控,Token 成本高

1.3 编排模式决策树

你的任务是什么样的?
│
├─ 有固定的分析→执行→验证流水线
│   → Sequential Pipeline
│   例:文档审核(解析→规则检查→人工复核→归档)
│
├─ 任务类型多变,但子任务种类可预定义
│   → Supervisor/Orchestrator
│   例:研发助手(有时调研、有时写代码、有时审查,由一个中心调度)
│
├─ Agent 需要在对话中自动分流到不同专家
│   → Swarm/Handoff
│   例:智能客服(售前→销售Agent,售后→支持Agent,无需用户手动选择)
│
└─ 需要多方对抗或评审来提升输出质量
    → Adversarial/Review
    例:安全审计(一个Agent生成代码,另一个独立审查漏洞)

工程建议: 四种模式中,Supervisor 模式在灵活性与可控性之间找到了最佳平衡点,也是 LangGraph 和 CrewAI Enterprise 官方推荐的首选 Multi-Agent 模式。本章后续代码也将基于 Supervisor 模式展开。


二、任务拆解:从自然语言到 Agent 可执行计划

2.1 拆解的两大范式

任务拆解是 Multi-Agent 编排的第一个核心技术决策。两种主流范式:

范式 原理 优点 缺点
层级规划 (Plan-then-Execute) 先让 Planner Agent 生成完整任务树,再逐节点分配执行 全局最优、可追踪 环境变化后计划失效
ReAct 动态补全 每步执行后根据观察结果决定下一步,边做边拆 自适应、容错 可能陷入局部最优或循环

工程最佳实践:Plan-then-Execute 做骨架 + ReAct 做纠偏。 即先用 Planner 生成结构化任务 DAG,执行过程中每个 Worker 有权在遇到异常时触发 RePlan 请求,由 Supervisor 局部重排。

2.2 任务拆解的输出格式

拆解结果必须结构化,不能让 Agent 之间用自然语言传递任务——这是 2025 年踩得最多的坑。推荐的标准化格式:

from typing import List, Literal
from pydantic import BaseModel, Field

class Subtask(BaseModel):
    """单个子任务"""
    id: str                           # e.g. "S1", "S2.1"
    description: str                  # 人类可读描述
    agent_type: Literal["research", "code", "review", "report"]
    dependencies: List[str] = Field(default_factory=list)
    tools_required: List[str] = Field(default_factory=list)
    expected_output_schema: str       # 输出 Schema 的类型名

class TaskPlan(BaseModel):
    """Planner 输出的完整计划"""
    goal: str
    subtasks: List[Subtask]
    execution_order: List[List[str]] = Field(default_factory=list)
    fallback_strategy: str            # 全局失败时的降级策略

2.3 Planner Agent 的 Prompt 设计

PLANNER_SYSTEM_PROMPT = """你是一个任务规划专家。给定用户目标,你需要将目标拆解为可分配给专业 Agent 的子任务。

拆解规则:
1. 每个子任务必须由一个特定类型的 Agent 独立完成
2. 子任务间的依赖关系必须形成 DAG(无环)
3. 如果某个子任务可能失败,必须指定 fallback 方案
4. 输出必须严格遵循 TaskPlan Schema

可用 Agent 类型及其能力:
- research: 信息检索、网页搜索、文档分析
- code: 编写/修改代码、运行测试、调试
- review: 代码审查、安全检查、质量评估
- report: 汇总输出、格式化文档、生成报告

可用工具:
- web_search(query): 网络搜索
- read_file(path): 读取文件
- run_test(path): 运行测试
- code_review(file): 代码审查
"""

def plan_task(user_goal: str, llm) -> TaskPlan:
    """让 Planner Agent 生成结构化任务计划。

    注意:llm 需支持 structured output(如 GPT-4o 的 response_format
    或 LangChain 的 with_structured_output),否则需在 Prompt 中要求
    输出 JSON 后手动解析。
    """
    structured_llm = llm.with_structured_output(TaskPlan)
    return structured_llm.invoke(
        f"{PLANNER_SYSTEM_PROMPT}\n\n请为以下目标生成任务计划:{user_goal}"
    )

2.4 执行过程中的动态调整

以下是一个概念示例——展示 Supervisor 如何在执行阶段追踪计划进度并触发重规划。注意这使用了 Pydantic 模型以便附加方法;在实际 LangGraph 实现中(见 §6),状态通常使用 TypedDict 以获得更好的序列化支持。

from typing import List
from pydantic import BaseModel, Field

class SupervisorState(BaseModel):
    plan: TaskPlan
    completed: List[str] = Field(default_factory=list)
    results: dict = Field(default_factory=dict)
    retry_counts: dict = Field(default_factory=dict)

    class Config:
        # 允许 int 作为 key 的 dict(规避 Pydantic 的类型强制转换)
        arbitrary_types_allowed = True

    MAX_RETRIES: int = 2

    def should_replan(self, failed_id: str) -> bool:
        """某个子任务在重试上限后仍失败 → 触发重规划"""
        self.retry_counts[failed_id] = self.retry_counts.get(failed_id, 0) + 1
        return self.retry_counts[failed_id] > self.MAX_RETRIES

    def get_next_batch(self) -> List[Subtask]:
        """返回依赖已满足且未完成的下一批子任务"""
        ready = []
        for st in self.plan.subtasks:
            if st.id in self.completed:
                continue
            if all(dep in self.completed for dep in st.dependencies):
                ready.append(st)
        return ready

三、工具调用:注册、Schema 与容错设计

3.1 工具注册中心

Multi-Agent 系统中每个 Agent 只能访问自己需要的工具——这叫最小权限原则(Least Privilege),是防止工具滥用和降低选择错误率的关键:

from typing import Callable, Any
from pydantic import BaseModel, Field

class ToolDef(BaseModel):
    """工具定义"""
    name: str
    description: str           # Agent 用这段描述决定何时调用
    parameters: dict           # JSON Schema
    func: Callable = Field(exclude=True)  # 实际函数,不序列化
    require_approval: bool = False       # 是否需要人工审批
    timeout_sec: float = 30.0

class ToolRegistry:
    """工具注册中心 — 每个 Agent 有独立的工具集"""

    def __init__(self):
        self._tools: dict[str, ToolDef] = {}

    def register(self, tool: ToolDef):
        self._tools[tool.name] = tool

    def get_for_agent(self, tool_names: List[str]) -> List[dict]:
        """返回 Agent 可见的工具列表(OpenAI Function Calling 格式)"""
        return [
            {
                "type": "function",
                "function": {
                    "name": self._tools[name].name,
                    "description": self._tools[name].description,
                    "parameters": self._tools[name].parameters,
                }
            }
            for name in tool_names if name in self._tools
        ]

    def execute(self, name: str, args: dict) -> Any:
        """执行工具调用,带异常捕获。

        注意:这里未实现真正的超时中断(signal 模块不可跨平台),
        生产环境建议使用 concurrent.futures.ThreadPoolExecutor + timeout 参数。
        """
        tool = self._tools.get(name)
        if not tool:
            return {"__tool_error": f"Tool '{name}' not found"}
        try:
            result = tool.func(**args)
            return {"__tool_result": result}
        except Exception as e:
            return {"__tool_error": str(e), "tool": name}


# ── 注册示例 ─────────────────────────────
registry = ToolRegistry()

registry.register(ToolDef(
    name="web_search",
    description="搜索互联网获取最新信息。输入 query 字符串。",
    parameters={
        "type": "object",
        "properties": {
            "query": {"type": "string", "description": "搜索关键词"}
        },
        "required": ["query"]
    },
    func=lambda query: f"搜索结果 for '{query}': ...",
))

registry.register(ToolDef(
    name="run_test",
    description="运行测试套件。输入 test_path 指定测试文件。",
    parameters={
        "type": "object",
        "properties": {
            "test_path": {"type": "string"}
        },
        "required": ["test_path"]
    },
    func=lambda test_path: f"测试 {test_path} 通过: 10/10",
    timeout_sec=120.0,
))

3.2 工具调用的错误恢复闭环

工具调用不是"调用→成功"那么简单,必须设计三级容错:

Level 1: 单工具重试(网络抖动、API 限流)
    → 指数退避,最多 3 次
Level 2: 工具替换(这个 API 挂了,换备用的)
    → 工具注册时声明 fallback_tool
Level 3: 任务降级(工具链整体不可用,用简化方案)
    → 触发 Supervisor 的 fallback_strategy
from typing import Any
import asyncio

class ResilientToolExecutor:
    """带三级容错的工具执行器"""

    def __init__(self, registry: ToolRegistry):
        self.registry = registry

    async def execute_with_retry(
        self, name: str, args: dict,
        max_retries: int = 3,
        base_delay: float = 1.0,
    ) -> Any:
        """Level 1: 指数退避重试"""
        for attempt in range(max_retries):
            result = self.registry.execute(name, args)
            if "__tool_error" not in result:
                return result["__tool_result"]
            if attempt < max_retries - 1:
                delay = base_delay * (2 ** attempt)
                await asyncio.sleep(delay)
        return {"__tool_error": f"Tool '{name}' failed after {max_retries} retries", "detail": result.get("__tool_error", "")}

    async def execute_with_fallback(
        self, name: str, args: dict,
        fallback_name: str = None,
    ) -> Any:
        """Level 2: 主工具失败后尝试备选"""
        result = await self.execute_with_retry(name, args)
        if "__tool_error" not in result:
            return result
        if fallback_name:
            return await self.execute_with_retry(fallback_name, args)
        return result

3.3 多工具并行调用

当多个工具调用之间无依赖时(如同时搜索多个数据源),并行执行可大幅压缩端到端延迟:

async def parallel_tool_calls(
    executor: ResilientToolExecutor,
    calls: list[tuple[str, dict]],   # [(tool_name, args), ...]
) -> dict[str, Any]:
    """并行执行多个独立的工具调用。

    返回格式:{tool_name: result_value},失败时值为错误信息字符串。
    注意:仅对无依赖关系的工具调用使用并行,有依赖的调用必须串行。
    """
    tasks = [
        executor.execute_with_retry(name, args)
        for name, args in calls
    ]
    results = await asyncio.gather(*tasks, return_exceptions=True)

    output = {}
    for i, (name, _) in enumerate(calls):
        r = results[i]
        if isinstance(r, Exception):
            output[name] = f"[并行调用异常] {r}"
        elif isinstance(r, dict) and "__tool_error" in r:
            output[name] = f"[工具错误] {r['__tool_error']}"
        else:
            output[name] = r
    return output

四、记忆模块:三层架构设计与实现

4.1 为什么 Agent 记忆不是 RAG

传统的"把历史对话全塞进向量数据库然后 top-k 检索"方案在 Multi-Agent 系统中会迅速崩塌——因为不同 Agent 需要不同类型的上下文:

  • Researcher Agent 需要"上周做过的类似调研",按任务语义检索
  • Coder Agent 需要"上次改这个文件时的上下文",按文件路径检索
  • Reviewer Agent 需要"项目整体的安全策略",按实体关系检索

一个统一的向量索引无法同时高效服务这三种查询模式。认知科学将人类记忆分为工作记忆、情节记忆和语义记忆三个层次,现代 Agent 记忆架构直接借鉴了这套体系。

层次 对应概念 存储 时效 典型容量 访问延迟
Working Memory 短期工作记忆 LLM Context Window 单次会话内 受限于上下文窗口 0ms(已在 prompt 中)
Episodic Memory 情节/事件记忆 向量数据库 + KV Store 跨会话持久 百万级事件 ~10–50ms
Semantic Memory 语义/知识记忆 向量数据库 + 知识图谱 长期稳定 TB 级 ~50ms+

4.2 信息流转路径

用户输入
    │
    ▼
┌─────────────────┐
│  感知层          │ ← 从 Episodic DB 检索相关历史
│  (注入 Context)  │ ← 从 Semantic KB 检索领域知识
└────────┬────────┘
         ▼
┌─────────────────┐
│  Working Memory  │ ← LLM 推理(Prompt 中包含检索结果)
└────────┬────────┘
         ▼
┌─────────────────┐
│  LLM 输出        │ ──→ 会话结束
└────────┬────────┘
         ▼
    新的 Episodic 记录(压缩→向量化→写入)
         │
         └──→ 定期升华至 Semantic Memory(高频稳定知识)

4.3 Working Memory:Context Window 的精打细算

Working Memory 的管理是 Agent 系统最大的隐性成本。以 GPT-4o 128K 窗口为例,扣除 System Prompt、工具 Schema、历史对话后,实际可用空间通常只剩 20%–40%。

推荐的 Prompt 结构(Token 预算分配):

[System Prompt: ~500 tokens]       定义 Agent 角色 + 行为约束
[Tool Schemas: ~200/工具]          每个工具的 JSON Schema
[Retrieved Episodic: ~2000 tokens]  向量检索到的相关历史(top-3)
[Retrieved Semantic: ~1500 tokens]  知识图谱检索的领域知识(top-3)
[Conversation History: 滑动窗口]    最近 N 轮对话,超出用摘要压缩
[Current Task: 剩余空间]            当前要执行的指令

滑动窗口 + 摘要压缩的混合策略:

class WorkingMemoryManager:
    """管理 Context Window 中的会话历史"""

    def __init__(self, max_history_tokens: int = 8000):
        self.max_tokens = max_history_tokens
        self.messages: list = []
        self.summary: str = ""  # 超出窗口的旧对话摘要

    def add(self, role: str, content: str, token_count: int):
        self.messages.append({
            "role": role,
            "content": content,
            "tokens": token_count,
        })
        self._compact()

    def _compact(self):
        """当总 token 超限时,将最早的消息移入摘要。

        token_count 为调用方传入的估算值(如 len(content)//4),
        精确值应使用 tiktoken 等分词器计算。
        """
        # 粗略 token 估算:英文 ~4 chars/token,中文 ~1.5 chars/token
        summary_tokens = len(self.summary) // 4
        total = sum(m["tokens"] for m in self.messages) + summary_tokens
        while total > self.max_tokens and len(self.messages) > 2:
            oldest = self.messages.pop(0)
            self.summary = self._update_summary(self.summary, oldest)
            summary_tokens = len(self.summary) // 4
            total = sum(m["tokens"] for m in self.messages) + summary_tokens

    def _update_summary(self, current: str, msg: dict) -> str:
        """将旧消息追加到摘要中。生产环境可改为调用轻量 LLM 做摘要合并。"""
        snippet = msg["content"][:200]  # 字符截断(非精确 token 截断)
        if not current:
            return f"[摘要] {msg['role']}: {snippet}"
        return current + f" | {msg['role']}: {snippet}"

    def get_context(self) -> str:
        """组装最终注入 LLM 的上下文"""
        parts = []
        if self.summary:
            parts.append(f"## 历史摘要\n{self.summary}")
        parts.append("## 最近对话")
        parts.extend(f"{m['role']}: {m['content']}" for m in self.messages)
        return "\n\n".join(parts)

4.4 Episodic Memory:向量检索 + 上下文扩展

Episodic Memory 存储的是"谁在什么时候做了什么"的事件记录。关键在于检索时不仅召回 top-k 匹配片段,还要扩展相邻上下文——因为语义上相关的信息往往分散在时间上相邻的多条记录中。MemMachine(MemVerge, 2026)的研究表明,上下文扩展可将召回准确率提升 4–8 个百分点。

import chromadb
from chromadb.utils import embedding_functions
import time

class EpisodicMemory:
    """基于 ChromaDB 的情节记忆"""

    def __init__(self, collection_name: str = "agent_episodes"):
        self.client = chromadb.PersistentClient(path="./memory_db")
        self.ef = embedding_functions.SentenceTransformerEmbeddingFunction(
            model_name="all-MiniLM-L6-v2"
        )
        self.collection = self.client.get_or_create_collection(
            name=collection_name,
            embedding_function=self.ef,
        )

    def store(self, agent_id: str, session_id: str,
              event_type: str, content: str, metadata: dict = None):
        """存储单条事件"""
        self.collection.add(
            documents=[content],
            metadatas=[{
                "agent_id": agent_id,
                "session_id": session_id,
                "event_type": event_type,
                "timestamp": time.time(),
                **(metadata or {}),
            }],
            ids=[f"{session_id}_{int(time.time()*1000)}"],
        )

    def retrieve(
        self, query: str,
        n_results: int = 3,
        context_expand: int = 1,
    ) -> list[dict]:
        """检索 + 上下文扩展。

        注意:ChromaDB 的 get() 不保证返回顺序。
        生产环境应使用带时间戳排序的专用查询,或迁移到支持
        时序过滤的向量数据库(如 Milvus、Weaviate)。
        """
        results = self.collection.query(
            query_texts=[query],
            n_results=n_results,
        )
        expanded = []
        for doc_id in results["ids"][0]:
            # 提取 session_id,取该会话中相邻的事件
            session_id = doc_id.split("_")[0]
            neighbors = self.collection.get(
                where={"session_id": session_id},
                limit=context_expand * 2 + 1,
            )
            expanded.extend(
                {"content": d, "meta": m}
                for d, m in zip(neighbors["documents"], neighbors["metadatas"])
            )
        return expanded

4.5 Semantic Memory:知识图谱 + 向量混合检索

对于"用户偏好 Python"、"项目使用 FastAPI"这类稳定的结构化知识,向量检索不够——需要知识图谱来维护实体间的关系。

# Semantic Memory:键值 + 三元组的轻量实现。
# 注意:此为演示用途的内存实现,不持久化。
# 生产环境建议使用 Neo4j 或带图扩展的 PostgreSQL(如 AGE)做知识图谱存储。
class SemanticMemory:
    """语义记忆:存储 Agent 学习到的长期知识"""

    def __init__(self):
        self.facts: dict[str, list[str]] = {}  # {entity: [fact1, fact2]}
        self.relations: list[tuple] = []        # [(entity1, relation, entity2)]

    def learn(self, entity: str, fact: str):
        if entity not in self.facts:
            self.facts[entity] = []
        self.facts[entity].append(fact)

    def learn_relation(self, e1: str, rel: str, e2: str):
        self.relations.append((e1, rel, e2))

    def query(self, entity: str) -> dict:
        """查询实体的所有已知事实和关系"""
        return {
            "facts": self.facts.get(entity, []),
            "relations": [
                (s, r, o) for s, r, o in self.relations
                if s == entity or o == entity
            ],
        }

    def to_context(self, entities: List[str]) -> str:
        """将知识图谱片段转为 LLM 可理解的文本"""
        lines = []
        for e in entities:
            info = self.query(e)
            if info["facts"]:
                lines.append(f"关于 {e} 的已知事实:")
                for f in info["facts"][-5:]:  # 取最近 5 条
                    lines.append(f"  - {f}")
            if info["relations"]:
                lines.append(f"{e} 的关联:")
                for s, r, o in info["relations"][-5:]:
                    lines.append(f"  - {s} --{r}--> {o}")
        return "\n".join(lines) if lines else ""

五、框架选型速查

2026 年三大 Multi-Agent 框架对比:

维度 LangGraph CrewAI AutoGen (AG2)
编排模型 有状态有向图(StateGraph) 角色-任务-流程声明式 多 Agent 对话
状态管理 内置 Checkpoint + State 持久化 弱(主要靠 Task 输入输出) 弱(靠对话历史)
流程控制 条件分支、循环、并行、人工审批节点 Sequential / Hierarchical 对话式协商
调试能力 LangSmith 时间旅行调试 日志输出 对话记录
MCP 支持 原生支持 部分支持 部分支持
学习曲线 陡(需要理解图编程) 平(声明式 YAML/Python)
适合场景 需要精确流程控制的生产系统 快速原型、角色分工明确的场景 需要多 Agent 辩论/推理的场景
开源 MIT Apache 2.0 MIT (原微软)

选型建议: Supervisor 模式用 LangGraph(精确控制),Sequential Pipeline 用 CrewAI(开发快),Adversarial 场景用 AutoGen(天然支持对话式对抗)。LangGraph + CrewAI 混合架构在生产中越来越常见——LangGraph 做流程引擎,CrewAI 在关键节点做角色协作。


六、完整实战代码:LangGraph Supervisor + 工具链 + 记忆

以下是一个 Multi-Agent Supervisor 系统的最小完整实现。注意:代码中的模块引用(agent_toolsagent_memory)对应前文定义的 ToolRegistryEpisodicMemory 等类——使用时需将它们组织到实际文件中。

"""
multi_agent_supervisor.py
Multi-Agent Supervisor 架构 — 可运行示例
依赖: pip install langgraph langchain-openai chromadb pydantic

注意:使用前需将前文定义的 ToolRegistry, EpisodicMemory 等类
放入 agent_tools.py 和 agent_memory.py 模块中。
"""
import json
from typing import TypedDict, List, Literal

from langgraph.graph import StateGraph, END
from langgraph.checkpoint.memory import MemorySaver
from langchain_openai import ChatOpenAI

from agent_tools import ToolRegistry, ResilientToolExecutor, parallel_tool_calls
from agent_memory import WorkingMemoryManager, EpisodicMemory, SemanticMemory


# ══════════════════════════════════════════════
# 1. 状态定义
# ══════════════════════════════════════════════

class SupervisorState(TypedDict):
    messages: List[dict]           # 全局消息队列
    next_agent: str                # Supervisor 路由到的下一个 Agent
    task: str                      # 当前任务
    session_id: str                # 会话 ID,用于记忆存储隔离
    plan: dict                     # 任务计划(子任务列表)
    results: dict                  # {agent_name: result_str}
    retry_counts: dict             # {agent_name: int}
    final_output: str              # 最终输出


# ══════════════════════════════════════════════
# 2. 初始化基础设施
# ══════════════════════════════════════════════

llm = ChatOpenAI(model="gpt-4o", temperature=0.1)

# 工具注册中心(详见 §3.1)
registry = ToolRegistry()
# 按需注册工具:registry.register(ToolDef(...))
executor = ResilientToolExecutor(registry)

# 记忆模块(详见 §4)
working_memory = WorkingMemoryManager(max_history_tokens=8000)
episodic = EpisodicMemory()
semantic = SemanticMemory()

MAX_RETRIES = 2


# ══════════════════════════════════════════════
# 3. Worker Agent 节点
# ══════════════════════════════════════════════

AGENT_CONFIGS = {
    "researcher": {
        "system": "你是一个信息检索专家。搜索、分析、总结信息。",
        "tools": ["web_search", "read_file"],
    },
    "coder": {
        "system": "你是一个软件工程师。编写、修改、测试代码。",
        "tools": ["read_file", "write_file", "run_test"],
    },
    "reviewer": {
        "system": "你是一个代码审查专家。检查安全性、性能、可维护性。",
        "tools": ["read_file", "code_review"],
    },
}


def create_worker_node(agent_type: str):
    """创建指定类型的 Worker Agent 节点。

    每个 Worker 节点:
    1. 从 Episodic + Semantic 记忆检索相关上下文
    2. 组装 Working Memory + 系统提示 + 任务
    3. 调用 LLM(支持工具调用)
    4. 将执行结果写回记忆和状态
    """
    config = AGENT_CONFIGS[agent_type]
    tools = registry.get_for_agent(config["tools"])

    def worker_node(state: SupervisorState) -> dict:
        task = state["task"]
        session_id = state.get("session_id", "default")

        # ── 从记忆中检索 ──
        relevant_episodes = episodic.retrieve(task, n_results=2)
        #  从 task 文本中提取可能涉及的实体名(简化:按空格分词取名词)
        entities = _guess_entities(task)
        semantic_context = semantic.to_context(entities)
        history_context = working_memory.get_context()

        # ── 组装 Prompt ──
        full_prompt = (
            f"{config['system']}\n\n"
            f"## 历史摘要\n{history_context}\n\n"
            f"## 相关事件(Episodic)\n{json.dumps(relevant_episodes, ensure_ascii=False)}\n\n"
            f"## 领域知识(Semantic)\n{semantic_context}\n\n"
            f"## 当前任务\n{task}"
        )

        # ── LLM 推理(带工具调用) ──
        response = llm.invoke(
            [
                {"role": "system", "content": full_prompt},
                {"role": "user", "content": task},
            ],
            tools=tools,
            tool_choice="auto",
        )

        # ── 处理工具调用 ──
        if hasattr(response, "tool_calls") and response.tool_calls:
            # LangChain 的 tool_calls 格式:{"name": ..., "args": ..., "id": ...}
            calls = [(tc["name"], tc["args"]) for tc in response.tool_calls]
            import asyncio
            try:
                loop = asyncio.get_event_loop()
                if loop.is_running():
                    # LangGraph 可能在已有事件循环中运行,不能用 asyncio.run()
                    import concurrent.futures
                    with concurrent.futures.ThreadPoolExecutor() as pool:
                        future = pool.submit(
                            lambda: asyncio.run(parallel_tool_calls(executor, calls))
                        )
                        tool_results = future.result(timeout=60)
                else:
                    tool_results = asyncio.run(parallel_tool_calls(executor, calls))
            except RuntimeError:
                tool_results = asyncio.run(parallel_tool_calls(executor, calls))

            working_memory.add(
                "tool", json.dumps(tool_results, ensure_ascii=False),
                token_count=len(str(tool_results)) // 4,
            )

        # ── 存储 Episodic Memory ──
        response_text = response.content if hasattr(response, "content") else str(response)
        episodic.store(
            agent_id=agent_type,
            session_id=session_id,
            event_type="task_execution",
            content=response_text,
        )

        # ── 返回状态更新(不原地修改,LangGraph 推荐模式)──
        new_results = {**state.get("results", {}), agent_type: response_text}
        return {"results": new_results}

    return worker_node


def _guess_entities(task: str) -> List[str]:
    """从任务描述中提取实体关键词(简化实现)。"""
    stopwords = {"的", "了", "是", "在", "和", "与", "或", "请", "分析", "修复", "执行"}
    words = task.replace(",", " ").replace("。", " ").split()
    return [w for w in words if len(w) > 1 and w not in stopwords][:5]


# ══════════════════════════════════════════════
# 4. Supervisor 节点(核心路由逻辑)
# ══════════════════════════════════════════════

SUPERVISOR_SYSTEM = """你是一个 Multi-Agent 系统的调度器。根据当前状态决定下一步行动。

可选动作:
1. assign: 将任务分配给一个 Worker Agent(researcher/coder/reviewer)
2. finish: 所有子任务已完成,结束流程
3. replan: 某个 Agent 重试超限但关键子任务未完成,触发重新规划

输出格式(只输出 JSON,不要额外文字):
{"action": "assign|finish|replan", "agent": "researcher|coder|reviewer|null", "reason": "..."}
"""


def supervisor_node(state: SupervisorState) -> dict:
    """Supervisor:检查进度 → 决策下一步 → 返回路由信息"""
    results = state.get("results", {})
    retries = state.get("retry_counts", {})

    # ── 检查是否有 Agent 超过重试上限 ──
    for agent, count in retries.items():
        if count > MAX_RETRIES:
            return {
                "final_output": json.dumps({
                    "status": "degraded",
                    "reason": f"Agent '{agent}' 超过重试上限 ({MAX_RETRIES})",
                    "completed": results,
                }, ensure_ascii=False),
                "next_agent": "FINISH",
            }

    # ── 调用 LLM 做路由决策 ──
    try:
        response = llm.invoke([
            {"role": "system", "content": SUPERVISOR_SYSTEM},
            {"role": "user", "content": json.dumps({
                "task": state.get("task", ""),
                "plan": state.get("plan", {}),
                "results": results,
                "retries": retries,
            }, ensure_ascii=False)},
        ])
        decision = json.loads(response.content)
    except (json.JSONDecodeError, Exception) as e:
        # LLM 输出非法 JSON → 安全降级:直接结束
        return {
            "final_output": json.dumps({
                "status": "supervisor_parse_error",
                "error": str(e),
                "completed": results,
            }, ensure_ascii=False),
            "next_agent": "FINISH",
        }

    action = decision.get("action", "finish")

    if action == "finish":
        return {
            "final_output": json.dumps(results, ensure_ascii=False),
            "next_agent": "FINISH",
        }
    elif action == "assign":
        return {
            "next_agent": decision.get("agent", "researcher"),
            "task": decision.get("reason", state.get("task", "")),
        }
    elif action == "replan":
        # 清空结果重新开始
        return {"results": {}, "next_agent": "planner"}

    # 兜底
    return {"next_agent": "FINISH"}


# ══════════════════════════════════════════════
# 5. 构建图 & 编译
# ══════════════════════════════════════════════

def build_supervisor_graph():
    builder = StateGraph(SupervisorState)

    # Worker 节点
    builder.add_node("supervisor", supervisor_node)
    builder.add_node("researcher", create_worker_node("researcher"))
    builder.add_node("coder", create_worker_node("coder"))
    builder.add_node("reviewer", create_worker_node("reviewer"))

    # Supervisor → Worker(条件路由)
    builder.add_conditional_edges(
        "supervisor",
        lambda s: s["next_agent"],
        {
            "researcher": "researcher",
            "coder": "coder",
            "reviewer": "reviewer",
            "FINISH": END,
        },
    )

    # Worker 执行完 → 回到 Supervisor
    for agent in ["researcher", "coder", "reviewer"]:
        builder.add_edge(agent, "supervisor")

    builder.set_entry_point("supervisor")

    # MemorySaver 仅用于开发/调试(内存中,重启即丢失)
    # 生产环境应使用 SqliteSaver 或 PostgresSaver
    return builder.compile(checkpointer=MemorySaver())


# ══════════════════════════════════════════════
# 6. 运行示例
# ══════════════════════════════════════════════

if __name__ == "__main__":
    graph = build_supervisor_graph()

    config = {
        "configurable": {"thread_id": "task-001"},
    }

    result = graph.invoke(
        {
            "task": "分析项目安全性并修复发现的高危漏洞",
            "messages": [],
            "session_id": "demo-001",
            "results": {},
            "retry_counts": {},
            "plan": {
                "subtasks": [
                    {"id": "S1", "agent_type": "researcher", "description": "扫描项目依赖"},
                    {"id": "S2", "agent_type": "coder", "description": "修复漏洞"},
                    {"id": "S3", "agent_type": "reviewer", "description": "审查修复"},
                ]
            },
            "final_output": "",
            "next_agent": "",
        },
        config=config,
    )

    output = result.get("final_output", "无结果")
    print(f"最终结果: {output}")
    print(f"Agent 执行记录: {result.get('results', {})}")

七、避坑指南

现象 根因 解法
Agent 间自然语言通信 解析失败、歧义导致错误级联 LLM 输出格式不稳定 必须用 Pydantic Schema 约束 Agent 间消息格式
单 Agent 挂 20+ 工具 工具选择频繁出错 LLM 在每个推理步骤中的注意力被过多工具 Schema 分散 每个 Agent 只给 3–8 个必需工具(最少权限);工具数超 8 就应拆 Agent
无限工具重试循环 Token 消耗飙升、任务卡死 工具失败后无上限重试 设置 MAX_RETRIES=2,超限后触发 replan/降级
Supervisor 单点故障 整个系统卡在 Supervisor 节点 Supervisor LLM 调用超时或报错 Supervisor 加超时(30s)+ 默认降级策略 + 多 Supervisor 热备
多 Agent 循环调用 A→B→A→B… 死循环 缺少全局调用图检测 在 Supervisor 中维护调用链,检测到 3 次以上循环后强制终止
Context Window 爆炸 对话越长 Agent 越"健忘" 历史对话 + 工具结果 + 检索内容超出窗口 滑动窗口 + 旧对话摘要压缩 + Token 预算监控
检索噪音淹没关键信息 向量检索返回无关历史 单纯 top-k 语义匹配不准 上下文扩展 + 时间衰减加权(最近的事件权重更高)
计划与现实脱节 Planner 第一次输出的计划执行到一半就过时 任务环境在变化 用 “Plan骨架 + ReAct纠偏” 模式,每个 Worker 有权申请 RePlan
人工审批节点缺失 危险操作(删库、发邮件)无人工确认 所有操作全自动执行 对高危工具加 require_approval=True,在 LangGraph 中插 interrupt_before
所有 Agent 用同一个记忆库 安全隔离失败,Agent A 读到 Agent B 的敏感数据 缺少记忆访问控制 每个 Agent 的记忆查询加 agent_id 过滤,敏感数据独立存储

八、趋势预判与总结

8.1 五大趋势

  1. Supervisor 模式成为生产默认选择。 2026 年头部 Agent 框架(LangGraph、CrewAI Enterprise)都将 Supervisor 作为推荐的多 Agent 编排模式——它在灵活性与可控性之间找到了最佳平衡。

  2. Agent 间通信从"自然语言"转向"结构化协议"。 OpenAI Swarm 的 handoff 机制、LangGraph 的 TypedState、微软 AutoGen 的结构化对话,都在推动 Agent 间通信的标准化。自由文本通信的可靠性问题已被广泛验证为生产噩梦。

  3. 记忆模块从"向量数据库"走向"向量 + 知识图谱 + 事件溯源"三位一体。 MemMachine(2026)和 Mem0 的流行表明:简单的 top-k 向量检索已不足以支撑跨 Session 的复杂查询。事件溯源(Event Sourcing)提供原始回放能力,知识图谱提供关系推理,向量检索提供语义入口——三者组合是 2026 年的标准范式。

  4. 人工审批节点从"可选"变为"合规刚需"。 EU AI Act 2026 年对高风险 AI 系统的监管要求逐步落地,多 Agent 系统涉及的关键决策点(如自动化审批、资金操作、内容发布)需要有 Human-in-the-Loop 审批机制。LangGraph 的 interrupt_before 和 CrewAI 的 human_input 回调正在成为标准配置。

  5. MCP(Model Context Protocol)统一工具生态。 Anthropic 推出的 MCP 协议正在被 LangGraph、CrewAI、AutoGen 全面集成——工具不再是每个框架独立定义,而是通过 MCP Server 统一暴露,Agent 通过标准协议发现和调用。

8.2 总览速查

┌──────────────────────────────────────────────────────────────────────┐
│          Multi-Agent 编排:任务拆解 + 工具调用 + 记忆模块 速查表        │
├────────────────────┬────────────────────┬────────────────────────────┤
│ 编排模式             │ Supervisor         │ 最灵活可控,生产首选         │
│                    │ Sequential         │ 固定流程,最简单             │
│                    │ Swarm/Handoff      │ 客服路由,动态转交           │
│                    │ Adversarial        │ 生成+批判,质量优先           │
├────────────────────┼────────────────────┼────────────────────────────┤
│ 任务拆解             │ Plan骨架 + ReAct纠偏│ 先规划再执行,允许动态调整    │
│ 任务间通信           │ Pydantic Schema    │ 必须结构化,不许自然语言      │
├────────────────────┼────────────────────┼────────────────────────────┤
│ 工具注册             │ ToolRegistry       │ 统一注册中心,最少权限        │
│ 工具容错             │ 重试→备选→降级     │ 三级容错闭环                 │
├────────────────────┼────────────────────┼────────────────────────────┤
│ 短时记忆 (Working)   │ 滑动窗口 + 摘要压缩 │ Token 预算监控               │
│ 情节记忆 (Episodic)  │ 向量检索 + 上下文扩展│ 优先保证召回,辅以精排   │
│ 语义记忆 (Semantic)  │ 知识图谱 + 向量混合  │ 关系推理能力                 │
├────────────────────┴────────────────────┴────────────────────────────┤
│ 核心法则:                                                            │
│ 1. 能用单 Agent 就不用 Multi-Agent — 复杂度是最贵的                    │
│ 2. Agent 间通信必须走结构化 Schema,禁止自由文本                       │
│ 3. 每个 Agent 只拿需要的工具(最少权限),超过 8 个就该拆 Agent         │
│ 4. 工具调用必须有三级容错(重试→备选→降级),不能假设一定成功          │
│ 5. Supervisor 必须设超时 + 默认降级策略,不能是单点                     │
│ 6. Context Window 是有限资源 — 滑动窗口 + 摘要压缩是必选项             │
│ 7. 高危操作(写数据库、发邮件、删文件)必须有 Human-in-the-Loop         │
└──────────────────────────────────────────────────────────────────────┘

参考来源:


摘要总结

多智能体编排的核心不是框架选择,而是三个工程基座的质量。任务拆解必须用 Plan骨架 + ReAct纠偏的混合模式——先规划全局 DAG 再允许动态重排,且 Agent 间通信必须走 Pydantic Schema 而非自然语言。工具调用遵循最少权限原则(每个 Agent 3–8 个工具)和三级容错闭环(重试→备选→降级)。记忆系统采用 Working/Episodic/Semantic 三层架构:短时记忆用滑动窗口 + 摘要压缩应对 Token 预算,情节记忆用向量检索 + 上下文扩展提升召回,语义记忆用知识图谱维护长期知识。2026 年的生产趋势是 LangGraph Supervisor 模式 + MCP 统一工具协议 + Human-in-the-Loop 合规审批——一句话:能用单 Agent 就别上 Multi-Agent,能结构化就别自由文本,能容错就别假设成功。

Logo

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

更多推荐