前言

做多智能体系统的人,多半都经历过这种场面:你起了一个 Planner Agent、一个 Coder Agent、一个 Reviewer Agent,让它们"互相讨论"完成任务。一开始看着挺像那么回事——它们你一句我一句,进度似乎在推进。但跑着跑着就开始出幺蛾子:约束莫名其妙没了,一个编造的 API 被另一个 Agent 当真接着用,两个 Agent 对"做完了"的理解根本不是一回事,对话轮次一多 token 直接爆掉。

问题不在于模型不够强,而在于让 Agent 用自然语言"对话"来协作,本身就是一种脆弱的协作方式。自然语言适合人和人沟通,却不适合作为 Agent 之间的工程契约。本文要讲的,就是用一套 A2A(Agent-to-Agent)结构化协议,加上 Orchestrator-Subagent 编排模式,把多智能体协作从"各说各话"变成"按契约交接"。

背景或问题

先看清楚多智能体自然语言对话式协作到底会在哪里翻车。下面四个场景都是真实工程里高频出现的。

翻车点 1:信息丢失——关键约束在转述中被稀释

Planner 把任务交给 Coder 时,原始约束是「用 PostgreSQL 13,预算 5000 元以内」。这个约束被 Coder 用自然语言复述给 Reviewer,Reviewer 又复述给自己当上下文。两轮转述之后,PostgreSQL 13 变成了"关系型数据库",预算约束干脆被省略。最终交付用了 MySQL 且超了预算。

自然语言是"有损压缩"。每经过一次转述,细节就会被取舍一次。Agent 之间没有强制的字段约束,谁都觉得"预算"这种细节"显然不重要",于是它就没了。

翻车点 2:幻觉传递——一个编造被多个 Agent 共同背书

Researcher 在搜集资料时幻觉编造了一个不存在的库 get_user_profile_v3()。Coder 拿来直接用,写得有模有样;Reviewer 看着"调用链很合理",也认为没问题。最后上线直接 ImportError。

最可怕的不是单个幻觉,而是幻觉会沿着对话链被反复"背书"。每多一个 Agent 采纳它,它看起来就越像事实。这和错误信息在人类群聊里越传越真是一个道理。

翻车点 3:协议不一致——对"完成"的定义各执一词

Coder 说"我做完了",它的意思是"函数写完了";Reviewer 听到"做完了",理解成"函数 + 测试 + 文档都齐了"。结果 Reviewer 开始挑刺,Coder 觉得自己被无理取闹,两个 Agent 互相踢皮球——要么死循环,要么 Reviewer 妥协提前放行,留下一堆没测试的代码。

各 Agent 对「任务完成」「需要补充信息」「我放弃」这些状态没有统一定义,协作语义就是散的。

翻车点 4:Token 爆炸——上下文随轮次指数膨胀

为了让每个 Agent 都"了解全局",每轮对话都把全部历史塞进上下文。三个 Agent 跑 8 轮,单次调用的上下文从 2k 涨到 30k token,成本指数上升;更糟的是,上下文一长还会触发 Lost in the Middle——模型对中间段信息注意力下降,质量反而变差。花了更多钱,得到了更差的结果。

在这里插入图片描述

这四个翻车点有一个共同根因:Agent 之间缺少一份机器可读、字段明确的协作契约。自然语言承诺不是工程契约——这个道理和「在 Prompt 里写请返回 JSON 不可靠」完全一样,只是被搬到了 Agent 之间。

核心思路

解决思路分两层。

第一层:通信契约——A2A 结构化协议。 Agent 之间不再发自然语言消息,而是交换结构化的"任务卡 / 结果卡"。任务卡规定"要做什么、输入是什么、约束是什么、超时多少";结果卡规定"做完了没、产出是什么、置信度多少、要不要重试"。状态机统一定义 pending / running / completed / failed / needs_input,所有 Agent 对状态的理解一致。

第二层:调度结构——Orchestrator-Subagent 编排模式。 引入一个编排器作为唯一的全局视角持有者,负责把大任务拆成子任务卡、调度子 Agent 执行、汇聚结果、裁决冲突。子 Agent 之间不直接对话,只和编排器交接卡片。

这两层合起来,把"自由对话"换成"结构化交接",从根上压制信息丢失、幻觉传递、协议不一致和 Token 爆炸。

实现步骤

1. 两种协作模式对比:P2P 对等 vs Orchestrator-Subagent

在写代码前,先想清楚用哪种协作拓扑。

P2P 对等协作:各 Agent 地位平等,互相直接通信,没有一个中心调度者。

  • 优点:灵活、去中心化,适合开放式探索、头脑风暴、多视角讨论。
  • 缺点:全局状态难追踪,容易陷入死循环,调试困难,成本和延迟不可控,结果冲突时无人裁决。
  • 适用场景:探索性、没有明确交付物的任务,比如"几个 Agent 从不同角度讨论一个方案"。

Orchestrator-Subagent 编排:一个编排器统一调度,子 Agent 只和编排器通信,彼此隔离。

  • 优点:全局可控、可观测、可并行、冲突可裁决、Token 可治理。
  • 缺点:编排器是单点和潜在瓶颈,编排逻辑本身复杂时可能成为脆弱点;灵活性不如 P2P。
  • 适用场景:有明确交付物、可拆解的工程任务,尤其生产环境。
    在这里插入图片描述

一句话选型:任务能拆、要交付、要上线,选 Orchestrator-Subagent;任务开放、要发散、不卡交付,再考虑 P2P。 本文聚焦前者。

2. A2A 协议:用 Pydantic 定义消息契约

A2A 协议的核心是三件套:任务卡(TaskCard)、结果卡(ResultCard)、状态机(TaskStatus)。

  • 任务卡是编排器发给子 Agent 的输入契约,字段包括 task_id、goal、inputs、constraints、timeout、parent_id。子 Agent 只看这一张卡。
  • 结果卡是子 Agent 回传给编排器的输出契约,字段包括 task_id、status、output、confidence、needs_retry、note。子 Agent 完成后只回传这张卡,而不是整段对话历史。
  • 状态机统一了所有 Agent 对任务阶段的认知,杜绝"做完了"语义歧义。

用 Pydantic v2 定义的好处是字段自带校验、类型明确、可序列化,本身就是一份活的契约文档。

3. 编排器的四项职责

Orchestrator 不是简单转发器,它承担四项核心职责:

  1. 任务拆解:把一个高层目标拆成若干有依赖关系的子任务卡,标注哪些能并行、哪些必须串行。
  2. 并行调度:对无依赖的子任务并发执行(线程池/异步队列),缩短整体延迟。
  3. 结果汇聚:把各子 Agent 的结果卡按依赖关系拼接成最终交付物。
  4. 冲突解决:当多个 Agent 对同一子任务给出冲突结果时,按策略裁决——例如取置信度最高者,置信度并列则触发重试或人工介入。

4. Token 爆炸治理:上下文隔离 + 摘要回传

这是 Orchestrator-Subagent 模式自带的"福利",但值得单独讲清。

上下文隔离:每个子 Agent 启动时上下文是干净的,只接收属于自己的那张任务卡,看不到全局历史、也看不到别的子 Agent 的对话。Researcher 不需要知道 Drafter 和 Reviewer 在聊什么,它只管把自己的资料点交出来。

摘要回传:子 Agent 完成后,不把内部多轮对话历史回传,而是只回传一张结构化的结果卡。子 Agent 内部可能调了 5 次 LLM、产生了几千 token 的中间过程,但回传给编排器的只有几十 token 的结构化产出。长任务下,上下文增长从指数级被压到近线性。

在这里插入图片描述

这两个机制直接对应翻车点 4:不再有"每轮把全部历史塞进去"的问题。

代码示例

下面给出一个完整、可运行的 Python 多智能体编排 Demo。场景是「生成一份技术调研报告」,拆成「搜集资料 → 撰写初稿 → 审校润色」三个串行子任务。为了让你复制就能跑通,Demo 用 mock 的 LLM 调用,不依赖任何 API Key;真实项目里把 _mock_llm 换成实际模型调用即可。

环境说明

  • Python 3.10+
  • pydantic v2(pip install pydantic

完整代码

# -*- coding: utf-8 -*-
"""
多智能体 A2A 协议 + Orchestrator-Subagent 编排 Demo
场景:技术调研报告生成(搜集资料 -> 撰写初稿 -> 审校润色)
环境:Python 3.10+,pydantic v2
依赖:pip install pydantic
"""
from __future__ import annotations

import uuid
from enum import Enum
from typing import Any, Callable, Optional

from pydantic import BaseModel, Field


# ============================================================
# 一、A2A 协议:Agent 间结构化消息契约
# ============================================================

class TaskStatus(str, Enum):
    """任务状态机:统一定义 Agent 对'任务处于哪个阶段'的理解。
    所有 Agent 共享同一套状态语义,杜绝'做完了'这种歧义。"""
    PENDING = "pending"          # 已创建,等待调度
    RUNNING = "running"          # 子 Agent 正在执行
    COMPLETED = "completed"      # 成功完成,产出可用
    FAILED = "failed"            # 失败(可按重试策略再调度)
    NEEDS_INPUT = "needs_input"  # 需要补充信息,返回给编排器裁决


class TaskCard(BaseModel):
    """任务卡:编排器 -> 子 Agent 的输入契约。
    子 Agent 只看这一张卡,不需要看全局历史(上下文隔离)。"""
    task_id: str = Field(
        default_factory=lambda: uuid.uuid4().hex[:8],
        description="任务唯一 ID",
    )
    goal: str = Field(description="任务目标,一句话讲清交付什么")
    inputs: dict[str, Any] = Field(
        default_factory=dict, description="结构化输入参数"
    )
    constraints: list[str] = Field(
        default_factory=list, description="硬约束,如字数/格式/禁用项"
    )
    timeout_sec: float = Field(default=30.0, description="超时上限")
    parent_id: Optional[str] = Field(
        default=None, description="所属父任务,便于追溯"
    )


class ResultCard(BaseModel):
    """结果卡:子 Agent -> 编排器的输出契约。
    子 Agent 完成后只回传这张卡,而不是整段对话历史(摘要回传)。"""
    task_id: str
    status: TaskStatus
    output: Any = Field(default=None, description="结构化产出")
    confidence: float = Field(
        ge=0.0, le=1.0, default=1.0, description="置信度,供编排器裁决"
    )
    needs_retry: bool = Field(default=False, description="是否建议重试")
    note: str = Field(default="", description="简短说明 / 遗留问题")


# ============================================================
# 二、子 Agent 基类
# ============================================================

class Subagent:
    """子 Agent 基类:上下文隔离的最小执行单元。
    核心契约:run() 只接收 TaskCard,只返回 ResultCard。"""
    name: str = "base"

    def __init__(self, llm: Optional[Callable[[str], str]] = None):
        # llm 可注入;默认用 mock,避免读者必须有 API Key 才能跑通
        self.llm = llm or self._mock_llm

    @staticmethod
    def _mock_llm(prompt: str) -> str:
        """伪 LLM:演示用,真实项目里替换成实际模型调用"""
        return f"[mock-llm] 已处理:{prompt[:40]}..."

    def run(self, card: TaskCard) -> ResultCard:
        """基类负责状态机包装与异常兜底,子类只实现 execute"""
        try:
            return self.execute(card)
        except Exception as e:
            # 任何异常都转成结构化的 FAILED 结果卡,而不是抛崩溃
            # 这本身也是对'幻觉/异常传递'的防护
            return ResultCard(
                task_id=card.task_id,
                status=TaskStatus.FAILED,
                needs_retry=True,
                note=f"执行异常:{e}",
            )

    def execute(self, card: TaskCard) -> ResultCard:
        raise NotImplementedError("子类必须实现 execute")


# ============================================================
# 三、三个具体子 Agent
# ============================================================

class Researcher(Subagent):
    """搜集资料:根据主题产出要点清单"""
    name = "researcher"

    def execute(self, card: TaskCard) -> ResultCard:
        topic = card.inputs.get("topic", "未指定主题")
        # 子 Agent 内部可多次调用 self.llm,但只回传结构化结果卡
        _ = self.llm(f"请围绕 {topic} 搜集要点")  # 演示内部调用
        bullets = [
            f"{topic} 的核心原理与关键机制",
            f"{topic} 的典型应用场景与适用边界",
            f"{topic} 的已知局限与常见踩坑",
        ]
        return ResultCard(
            task_id=card.task_id,
            status=TaskStatus.COMPLETED,
            output={"bullets": bullets},
            confidence=0.85,
        )


class Drafter(Subagent):
    """撰写初稿:把资料要点组织成正文"""
    name = "drafter"

    def execute(self, card: TaskCard) -> ResultCard:
        bullets = card.inputs.get("bullets", [])
        if not bullets:
            # 缺输入时主动声明 NEEDS_INPUT,而不是瞎编——对抗幻觉传递
            return ResultCard(
                task_id=card.task_id,
                status=TaskStatus.NEEDS_INPUT,
                note="缺少资料要点,请先由 Researcher 产出 bullets",
            )
        topic = card.inputs.get("topic", "技术主题")
        body = "\n".join(f"- {b}" for b in bullets)
        draft = f"## {topic} 调研报告(初稿)\n{body}"
        return ResultCard(
            task_id=card.task_id,
            status=TaskStatus.COMPLETED,
            output={"draft": draft},
            confidence=0.8,
        )


class Reviewer(Subagent):
    """审校润色:检查初稿并产出终稿"""
    name = "reviewer"

    def execute(self, card: TaskCard) -> ResultCard:
        draft = card.inputs.get("draft", "")
        if not draft:
            return ResultCard(
                task_id=card.task_id,
                status=TaskStatus.NEEDS_INPUT,
                note="缺少初稿,请先由 Drafter 产出 draft",
            )
        final = draft + "\n\n## 小结\n(已审校:修正错别字、补充小结与结论)"
        return ResultCard(
            task_id=card.task_id,
            status=TaskStatus.COMPLETED,
            output={"final_report": final},
            confidence=0.9,
        )


# ============================================================
# 四、Orchestrator 编排器
# ============================================================

class Orchestrator:
    """编排器:唯一持有全局视角的角色。
    职责:任务拆解、调度、结果汇聚、冲突解决。子 Agent 之间不直接通信。"""

    def __init__(self):
        self.results: dict[str, ResultCard] = {}  # task_id -> 结果卡
        self.log: list[str] = []

    def _log(self, msg: str) -> None:
        print(f"[orchestrator] {msg}")
        self.log.append(msg)

    # ---- 职责 1:任务拆解 ----
    def decompose(self, topic: str) -> list[TaskCard]:
        """把高层目标拆成带依赖关系的子任务卡"""
        self._log(f"收到目标:生成《{topic}》技术调研报告,开始拆解")
        return [
            TaskCard(
                goal="搜集资料,产出要点清单",
                inputs={"topic": topic},
                constraints=["要点不少于 3 条", "必须基于事实,不得编造"],
                parent_id="root",
            ),
            TaskCard(
                goal="基于资料撰写初稿",
                inputs={"topic": topic},  # bullets 待阶段 1 回填
                constraints=["结构清晰", "不少于 200 字"],
                parent_id="root",
            ),
            TaskCard(
                goal="审校润色,产出终稿",
                inputs={},  # draft 待阶段 2 回填
                constraints=["无明显错别字", "补充小结"],
                parent_id="root",
            ),
        ]

    # ---- 职责 2 & 3:调度 + 汇聚 ----
    def schedule(self, cards: list[TaskCard]) -> str:
        """调度执行:演示带依赖的串行流水线。
        真实工程中,无依赖的子任务可用线程池/异步队列并发执行。"""
        researcher, drafter, reviewer = Researcher(), Drafter(), Reviewer()

        # 阶段 1:搜集资料
        r1 = researcher.run(cards[0])
        self.results[r1.task_id] = r1
        self._log(f"阶段1 搜集资料 -> {r1.status.value} conf={r1.confidence}")
        if r1.status != TaskStatus.COMPLETED:
            return self._fail("搜集资料未完成,流水线中止")

        # 阶段 2:撰写初稿(依赖阶段 1 产出,体现结构化交接)
        cards[1].inputs["bullets"] = r1.output["bullets"]
        r2 = drafter.run(cards[1])
        self.results[r2.task_id] = r2
        self._log(f"阶段2 撰写初稿 -> {r2.status.value}")
        if r2.status != TaskStatus.COMPLETED:
            return self._fail("撰写初稿未完成/缺输入,流水线中止")

        # 阶段 3:审校润色(依赖阶段 2 产出)
        cards[2].inputs["draft"] = r2.output["draft"]
        r3 = reviewer.run(cards[2])
        self.results[r3.task_id] = r3
        self._log(f"阶段3 审校润色 -> {r3.status.value}")
        if r3.status != TaskStatus.COMPLETED:
            return self._fail("审校未完成,流水线中止")

        return self.aggregate([r1, r2, r3])

    def aggregate(self, results: list[ResultCard]) -> str:
        """结果汇聚:从各阶段产出中提取最终交付物"""
        final = next(
            r for r in results if r.output and "final_report" in r.output
        )
        report = final.output["final_report"]
        self._log(f"汇聚完成,终稿字数约 {len(report)}")
        return report

    # ---- 职责 4:冲突解决 ----
    def resolve_conflict(self, results: list[ResultCard]) -> ResultCard:
        """冲突解决:多个 Agent 对同一子任务给出冲突结果时的裁决。
        策略:取置信度最高者;置信度并列则标记重试交人工裁决。"""
        results = sorted(results, key=lambda r: r.confidence, reverse=True)
        top = results[0]
        if len(results) > 1 and results[1].confidence == top.confidence:
            self._log("置信度并列,标记 needs_retry 交人工裁决")
            top.needs_retry = True
        return top

    def _fail(self, reason: str) -> str:
        self._log(f"FAILED:{reason}")
        return f"[FAILED] {reason}"


# ============================================================
# 五、运行入口
# ============================================================

if __name__ == "__main__":
    orch = Orchestrator()
    cards = orch.decompose(topic="多智能体协作")
    report = orch.schedule(cards)
    print("\n===== 最终报告 =====")
    print(report)

把串行升级成并发

上面的 Demo 是串行流水线,因为三个子任务有依赖。但当任务图里出现多个无依赖的子任务时(比如同时调研三个不同子主题),应该并发执行。把 schedule 里的串行调用换成线程池即可:

import concurrent.futures

def schedule_parallel(self, cards: list[TaskCard]) -> dict[str, ResultCard]:
    """并发调度无依赖的子任务,缩短整体延迟"""
    with concurrent.futures.ThreadPoolExecutor(max_workers=4) as pool:
        # 每个 Subagent 实例只接收自己的 TaskCard,天然互不干扰
        future_map = {
            pool.submit(self._agent_for(c).run, c): c.task_id for c in cards
        }
        for fut in concurrent.futures.as_completed(future_map):
            tid = future_map[fut]
            self.results[tid] = fut.result()
    return self.results

def _agent_for(self, card: TaskCard) -> Subagent:
    """根据任务卡 goal 路由到对应子 Agent(简化版路由)"""
    g = card.goal
    if "搜集" in g:
        return Researcher()
    if "初稿" in g:
        return Drafter()
    if "审校" in g:
        return Reviewer()
    raise ValueError(f"没有匹配的子 Agent:{g}")

并发之所以安全,正是因为每个子 Agent 上下文隔离、只吃自己的任务卡,没有共享可变状态。

运行结果或效果说明

直接 python demo.py 运行,输出大致如下(mock 数据):

[orchestrator] 收到目标:生成《多智能体协作》技术调研报告,开始拆解
[orchestrator] 阶段1 搜集资料 -> completed conf=0.85
[orchestrator] 阶段2 撰写初稿 -> completed
[orchestrator] 阶段3 审校润色 -> completed
[orchestrator] 汇聚完成,终稿字数约 180

===== 最终报告 =====
## 多智能体协作 调研报告(初稿)
- 多智能体协作 的核心原理与关键机制
- 多智能体协作 的典型应用场景与适用边界
- 多智能体协作 的已知局限与常见踩坑

## 小结
(已审校:修正错别字、补充小结与结论)

重点不在 mock 产出本身,而在过程结构:每一步交接的都是结构化卡片,编排器日志里能看到清晰的状态流转(completed / needs_input / failed),没有任何"自由对话"。把 _mock_llm 换成真实模型调用、把 mock 的 bullets 换成模型实际产出,这份骨架就能直接进项目。

常见问题与避坑

1. 编排器本身会不会变成新的瓶颈?
会。Orchestrator 是单点,它的拆解逻辑和冲突裁决逻辑一旦写错,全局都受影响。实践中要把编排逻辑也结构化、可测试:拆解规则用纯函数、裁决策略可配置、关键决策留审计日志。不要让编排器自己也变成一个"自由发挥的 Agent"。

2. 子 Agent 一定要无状态吗?
不绝对。子 Agent 内部可以有短期记忆(比如多次调用 LLM 的中间上下文),但这份记忆只在本次任务执行期间有效,任务结束、结果卡回传后就丢弃。关键是不要把内部记忆泄漏到 Agent 间通信里。

3. 结果卡的 output 用 Any 会不会太松?
会。Demo 里为了通用用了 Any,生产中应该给每类子任务定义具体的 output schema(比如 Researcher 用 ResearchOutput,Reviewer 用 ReviewOutput),让下游消费者能静态校验。这和单体 Agent 里的结构化输出是同一套思路。

4. 需要重试时怎么避免死循环?
给 needs_retry 加一个全局重试预算(如每个子任务最多重试 3 次),并在 ResultCard 里带上重试计数。超过预算就强制 FAILED 上报,由编排器决定降级还是终止整条流水线。状态机里必须有一条明确的"终止"出口。

5. 这套和 LangGraph / CrewAI / AutoGen 是什么关系?
这些框架都提供了多智能体编排能力,思想上是相通的——LangGraph 的图编排、CrewAI 的角色分工、AutoGen 的 GroupChat 都在解决类似问题。本文讲的是底层的协议与模式原理,理解之后用任何框架都能更清楚它在哪里帮你做了上下文隔离、在哪里做冲突裁决,而不是把它当黑盒。

总结

多智能体协作翻车的根因,往往是让 Agent 用自然语言对话来协作——信息会丢失、幻觉会被传递、协议会不一致、token 会爆炸。本文给出的解法是两层结构化:用 A2A 协议把 Agent 间通信从自然语言换成任务卡/结果卡/状态机,用 Orchestrator-Subagent 编排把协作从对等对话换成中心化调度。

落地时记住几个关键设计:任务卡和结果卡用 Pydantic 定义成强类型契约,让字段约束代替口头承诺;编排器专注拆解、调度、汇聚、裁决四件事,并且自身也要可测试、留审计;子 Agent 严格上下文隔离,完成后只回传结构化结果卡,把 token 增长从指数级压到近线性;重试要有预算,状态机要有终止出口,防止死循环。把这套骨架跑通,再往上接具体的模型和业务,多智能体系统才有可能从"看着好玩"走向"敢上线"。

Logo

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

更多推荐