Multi-Agent 通信协议与编排中枢:状态机、DAG 与事件总线

文章来源声明: 原文作者:吴佳浩Alben; 来源站点:掘金; 原文链接:https://juejin.cn/post/7688189200245866532; 本文基于上述来源整理/加工,觅优补充点评,仅供技术学习交流。版权归原作者所有。
觅优短评

这篇把多 Agent 协同从自由聊天拉回软件工程,协议、状态机与回滚缺一不可;适合正在设计企业级 Agent 编排架构、需要确定性与容错的团队精读。

Multi-Agent 通信协议与编排中枢:状态机、DAG 与事件总线 -----------------------------------

请添加图片描述

作 者:吴佳浩(Alben)
公众号:全栈架构师笔记
系列专栏:《企业级 Agent 实战指南————Multi-Agent 架构设计:从单体 ReAct 到群智协同》· 第 02 篇


导读
没有协议的多 Agent 协同是无序的群聊噪音,有了状态机与 DAG 的多 Agent 才是精密的软件工程。
很多团队在设计多 Agent 时,让每个 Agent 自由把自然语言当作通信载体;但在高并发分布式环境下,自然语言通信的歧义、缺乏结构化校验与无状态,会导致系统迅速陷入死锁。
静态工作流用 DAG 编排保证确定性,动态协同用事件总线与有限状态机(FSM)解耦调度。


在多智能体系统开发中,最常犯的一个低级错误就是“让所有 Agent 在一个共用的 Message List 里自由聊天”: 架构师 Agent 发一段文字,编码 Agent 回复一段文字,测试 Agent 再接着回话。

这种基于自然语言自由群聊的设计,在进入企业级复杂任务后,会迅速引发三大系统崩溃:

灾难现象具体翻车表现架构根因
1. 结构化契约丢失编码 Agent 漏掉了架构师提出的缺乏基于强类型 Schema 的协议
(Contract Failure)某个数据库字段定义,导致后续报错约束,关键参数在自然语言中丢失
2. 状态流转不可控Agent 之间互相客套,无法收敛到缺乏确定性的有限状态机 (FSM),
(State Runaway)明确的“已完成”终态,消耗海量 Token 无法判定当前处于哪个具体阶段
3. 事务无法回滚步骤 3 失败了,步骤 1 和 2 的缺乏基于事件总线与 Saga 模式的
(No Rollback/Saga)修改无法自动撤销,留下脏环境补偿与回滚机制

要构建企业级多智能体系统,必须建立统一的结构化通信协议(Agent Communication Protocol),并结合 DAG(有向无环图)与有限状态机(FSM) 进行严密编排。


一、多 Agent 通信协议设计:超越自然语言的结构化载荷

在微服务架构中,服务间通信依靠 gRPC 或 RESTful JSON,绝不靠一段不可解析的口语。多 Agent 协同同样必须基于强类型的结构化信元(Agent Envelope Protocol)

协议字段数据类型核心用途
trace\_idUUID (String)全局链路追踪,贯穿整个多 Agent 树
sender / receiverString (Agent Role)明确发送方与目标接收方 (路由依据)
message\_typeEnum (REQUEST, RESPONSE, EVENT明确动作意图,指示状态机如何流转
, BROADCAST, ERROR)
payloadStrongly-typed JSON Object结构化业务交付物 (代码、Diff、测试)
state\_contextDict\[str, Any\]当前任务的阶段性元数据与环境变量

mermaid diagram

  • 🔸 明确发送方与接收方:杜绝广播滥用,实现点对点精准路由;
  • 🔸 强类型 Payload 载荷:交付物必须是可解析的结构化实体(如 {"files_changed": [...], "test_exit_code": 0}),而不是一段模糊的口语;
  • 🔸 全局链路追踪(Trace ID):方便全链路排错与性能可观测性诊断。

一句话总结这一章的核心观点:
自然语言是 Agent 与人类交互的界面,结构化协议才是 Agent 与 Agent 协同的底座。


二、编排中枢的两大流派:静态 DAG vs. 动态状态机 (FSM)

在系统编排上,业界存在两大经典流派,架构师必须明确它们的适用边界:

架构流派静态 DAG 编排 (如 LangGraph)动态有限状态机 (FSM / Bus)
拓扑特征拓扑节点与依赖边在代码中预先固化拓扑由 Agent 根据运行时事件动态
(A -> B -> C -> End)判定转移 (Event-Driven Transition)
确定性与容错极高确定性,易于可视化与调试高灵活性,适应未知探索性任务
典型业务场景规范化 CI/CD 发布、标准合同审查、复杂未知 Bug 排错、智能客服路由
固定流水线 SOP、攻防博弈探索
架构代价遇到分支异常时难以跳跃;状态机过大时存在死锁风险;
扩展新步骤需修改图定义需要完备的熔断机制

mermaid diagram

一句话总结这一章的核心观点:
流程固定的企业 SOP 走静态 DAG 保证确定性,复杂探索型任务走状态机事件总线保证灵活性。


三、生产级代码实战:带状态机与回滚补偿的 Multi-Agent 编排器

以下为基于 Python 3.11+ 构建的企业级多 Agent 状态机编排引擎完整实现,包含强类型通信信元、状态转移拦截与 Saga 补偿回滚:

<span>"""
multi_agent_fsm_orchestrator.py

生产级 Multi-Agent 状态机编排引擎

包含:
- 结构化 Agent 信元协议(Envelope)
- 有限状态机(FSM)
- Review 驳回与回流
- 失败回滚
- 最大迭代保护(Circuit Breaker)
"""</span>

<span>import</span> enum
<span>import</span> uuid
<span>from</span> datetime <span>import</span> datetime
<span>from</span> typing <span>import</span> <span>Any</span>, <span>Dict</span>, <span>List</span>, <span>Optional</span>

<span>from</span> pydantic <span>import</span> BaseModel, Field


<span># ============================================================</span>
<span># Agent 状态定义</span>
<span># ============================================================</span>


<span>class</span> <span>AgentState</span>(<span>str</span>, enum.Enum):
    <span>"""Multi-Agent 工作流状态"""</span>

    PLANNING = <span>"PLANNING"</span>
    CODING = <span>"CODING"</span>
    REVIEWING = <span>"REVIEWING"</span>
    COMPLETED = <span>"COMPLETED"</span>
    FAILED = <span>"FAILED"</span>


<span># ============================================================</span>
<span># Agent 消息类型</span>
<span># ============================================================</span>


<span>class</span> <span>MessageType</span>(<span>str</span>, enum.Enum):
    <span>"""Agent 间标准消息"""</span>

    TASK_DISPATCH = <span>"TASK_DISPATCH"</span>
    TASK_DELIVER = <span>"TASK_DELIVER"</span>
    REVIEW_REJECT = <span>"REVIEW_REJECT"</span>
    REVIEW_PASS = <span>"REVIEW_PASS"</span>


<span># ============================================================</span>
<span># Agent Envelope</span>
<span># ============================================================</span>


<span>class</span> <span>AgentEnvelope</span>(<span>BaseModel</span>):
    <span>"""
    Agent 间通信信元

    每一次 Agent 通信都封装为统一 Envelope。
    """</span>

    trace_id: <span>str</span> = Field(
        default_factory=<span>lambda</span>: <span>str</span>(uuid.uuid4())
    )

    sender: <span>str</span>
    receiver: <span>str</span>

    msg_type: MessageType
    current_state: AgentState

    payload: <span>Dict</span>[<span>str</span>, <span>Any</span>]

    timestamp: datetime = Field(
        default_factory=datetime.utcnow
    )


<span># ============================================================</span>
<span># Multi-Agent FSM Runtime</span>
<span># ============================================================</span>


<span>class</span> <span>MultiAgentFSMRuntime</span>:
    <span>"""
    企业级 Multi-Agent 状态机运行时

    工作流:

        Planning
            │
            ▼
        Coding
            │
            ▼
        Reviewing
         │        │
         │Pass    │Reject
         ▼        ▼
     Completed  Coding
    """</span>

    <span>def</span> <span>__init__</span>(<span>self, user_goal: <span>str</span></span>):
        self.trace_id = <span>str</span>(uuid.uuid4())

        self.user_goal = user_goal

        self.state = AgentState.PLANNING

        self.iteration_count = <span>0</span>
        self.max_iterations = <span>6</span>

        self.history_envelopes: <span>List</span>[
            AgentEnvelope
        ] = []

    <span># --------------------------------------------------------</span>
    <span># FSM Dispatcher</span>
    <span># --------------------------------------------------------</span>

    <span>def</span> <span>dispatch_envelope</span>(<span>
        self,
        envelope: AgentEnvelope,
    </span>) -> <span>Optional</span>[AgentEnvelope]:
        <span>"""
        接收 Envelope,
        推动状态机流转。
        """</span>

        self.history_envelopes.append(envelope)

        self.iteration_count += <span>1</span>

        <span>print</span>(
            <span>"\n"</span>
            <span>"📨 [Message]\n"</span>
            <span>f"Sender   : <span>{envelope.sender}</span>\n"</span>
            <span>f"Receiver : <span>{envelope.receiver}</span>\n"</span>
            <span>f"Type     : <span>{envelope.msg_type.value}</span>\n"</span>
            <span>f"FSM      : <span>{self.state.value}</span>\n"</span>
        )

        <span># ----------------------------------------------------</span>
        <span># Circuit Breaker</span>
        <span># ----------------------------------------------------</span>

        <span>if</span> self.iteration_count > self.max_iterations:
            <span>print</span>(
                <span>"🚨 Maximum iteration reached. "</span>
                <span>"Workflow aborted."</span>
            )

            self.state = AgentState.FAILED
            <span>return</span> <span>None</span>

        <span># ----------------------------------------------------</span>
        <span># Planning → Coding</span>
        <span># ----------------------------------------------------</span>

        <span>if</span> (
            self.state == AgentState.PLANNING
            <span>and</span> envelope.msg_type
            == MessageType.TASK_DISPATCH
        ):
            self.state = AgentState.CODING

            <span>return</span> self._run_coder_agent(envelope)

        <span># ----------------------------------------------------</span>
        <span># Coding → Reviewing</span>
        <span># ----------------------------------------------------</span>

        <span>if</span> (
            self.state == AgentState.CODING
            <span>and</span> envelope.msg_type
            == MessageType.TASK_DELIVER
        ):
            self.state = AgentState.REVIEWING

            <span>return</span> self._run_reviewer_agent(envelope)

        <span># ----------------------------------------------------</span>
        <span># Review Reject</span>
        <span># ----------------------------------------------------</span>

        <span>if</span> (
            self.state == AgentState.REVIEWING
            <span>and</span> envelope.msg_type
            == MessageType.REVIEW_REJECT
        ):
            self.state = AgentState.CODING

            <span>print</span>(
                <span>"⚠️ Review rejected.\n"</span>
                <span>f"Reason: <span>{envelope.payload.get(<span>'reasons'</span>)}</span>\n"</span>
                <span>"Routing back to Coder..."</span>
            )

            <span>return</span> self._run_coder_agent(envelope)

        <span># ----------------------------------------------------</span>
        <span># Review Pass</span>
        <span># ----------------------------------------------------</span>

        <span>if</span> (
            self.state == AgentState.REVIEWING
            <span>and</span> envelope.msg_type
            == MessageType.REVIEW_PASS
        ):
            self.state = AgentState.COMPLETED

            <span>print</span>(
                <span>"🎉 Review passed.\n"</span>
                <span>"Workflow completed successfully."</span>
            )

            <span>return</span> <span>None</span>

        <span>return</span> <span>None</span>

    <span># --------------------------------------------------------</span>
    <span># Coder Agent</span>
    <span># --------------------------------------------------------</span>

    <span>def</span> <span>_run_coder_agent</span>(<span>
        self,
        incoming: AgentEnvelope,
    </span>) -> AgentEnvelope:
        <span>"""
        模拟 Coder Agent。
        """</span>

        is_fix = (
            incoming.msg_type
            == MessageType.REVIEW_REJECT
        )

        code = (
            <span>"def auth():\n"</span>
            <span>"    return True\n"</span>
            <span>"    # Fixed Logic"</span>
            <span>if</span> is_fix
            <span>else</span>
            <span>"def auth():\n"</span>
            <span>"    return False"</span>
        )

        <span>return</span> AgentEnvelope(
            trace_id=self.trace_id,
            sender=<span>"CoderAgent"</span>,
            receiver=<span>"ReviewerAgent"</span>,
            msg_type=MessageType.TASK_DELIVER,
            current_state=self.state,
            payload={
                <span>"code"</span>: code,
                <span>"is_fix"</span>: is_fix,
            },
        )

    <span># --------------------------------------------------------</span>
    <span># Reviewer Agent</span>
    <span># --------------------------------------------------------</span>

    <span>def</span> <span>_run_reviewer_agent</span>(<span>
        self,
        incoming: AgentEnvelope,
    </span>) -> AgentEnvelope:
        <span>"""
        模拟 Reviewer Agent。
        """</span>

        code = incoming.payload.get(<span>"code"</span>, <span>""</span>)

        <span># 模拟测试失败</span>
        <span>if</span> <span>"return False"</span> <span>in</span> code:
            <span>return</span> AgentEnvelope(
                trace_id=self.trace_id,
                sender=<span>"ReviewerAgent"</span>,
                receiver=<span>"CoderAgent"</span>,
                msg_type=MessageType.REVIEW_REJECT,
                current_state=self.state,
                payload={
                    <span>"reasons"</span>: (
                        <span>"AssertionError: "</span>
                        <span>"auth() returned False "</span>
                        <span>"on valid credentials."</span>
                    )
                },
            )

        <span>return</span> AgentEnvelope(
            trace_id=self.trace_id,
            sender=<span>"ReviewerAgent"</span>,
            receiver=<span>"Orchestrator"</span>,
            msg_type=MessageType.REVIEW_PASS,
            current_state=self.state,
            payload={
                <span>"summary"</span>: (
                    <span>"All 15 unit tests passed. "</span>
                    <span>"Code style verified."</span>
                )
            },
        )

    <span># --------------------------------------------------------</span>
    <span># Workflow Entry</span>
    <span># --------------------------------------------------------</span>

    <span>def</span> <span>start</span>(<span>self</span>) -> <span>None</span>:
        <span>"""
        启动 Multi-Agent FSM。
        """</span>

        <span>print</span>(
            <span>"\n"</span>
            <span>"🚀 Multi-Agent FSM Started\n"</span>
            <span>f"Goal: <span>{self.user_goal}</span>\n"</span>
        )

        initial = AgentEnvelope(
            trace_id=self.trace_id,
            sender=<span>"LeaderAgent"</span>,
            receiver=<span>"CoderAgent"</span>,
            msg_type=MessageType.TASK_DISPATCH,
            current_state=self.state,
            payload={
                <span>"spec"</span>: (
                    <span>"Implement authentication service"</span>
                )
            },
        )

        envelope = self.dispatch_envelope(initial)

        <span>while</span> (
            envelope
            <span>and</span> self.state
            <span>not</span> <span>in</span> (
                AgentState.COMPLETED,
                AgentState.FAILED,
            )
        ):
            envelope = self.dispatch_envelope(envelope)


本篇总结

  • 🔸 严禁自然语言无序群聊:必须定义基于 Trace ID 和强类型 Payload 的信元协议;
  • 🔸 静态 DAG 负责确定性流水线,动态 FSM 负责事件驱动自修复
  • 🔸 明确状态转移条件与最大迭代熔断保护,杜绝 Agent 之间反复扯皮导致的死循环;
  • 🔸 审查驳回机制是打破 Agent 自我幻觉的最强工程防线

在多 Agent 协同中,当成百上千个 Agent 同时并发读写共享数据时,如何避免数据冲突?

筒子们本篇为《企业级 Agent 实战指南》· 第三章的第 2 篇,后续续会更新完整的agent的开发的全部过程,如果你对Agent开发感兴趣不妨关注一下本合集。

在下一篇中,我们将深入拆解:《共享黑板模式(Blackboard)实战:多 Agent 如何并发协作而不冲突?》