多智能体系统(Multi-Agent System)火到今天,早就不停在概念层面了。真正上手做过多 Agent 协作项目的人,大概率都经历过同一个狼狈场景:Agent A 调 Agent B,Agent B 又去问 Agent C,消息满天飞,状态全靠日志猜,任务跑挂了你都说不清是哪一步、哪个节点、什么条件引起的。这套系统能不能撑住,核心不在这几个 LLM 模型有多聪明,而在通信和编排层能不能做到可控、可查、可恢复。这篇文章想拆的,就是我在这类系统里反复打磨的一个组合方案:用状态机管住行为,用 DAG 管住依赖,用事件总线管住通信,三者合起来才是真正能落地的编排中枢。适合正在搭多 Agent 服务的开发者、做自动化任务流编排的工程师看,也适合那些被 Agent 混乱协作折磨过的同学做一次复盘参考。
1. 多智能体系统为什么必须有一个编排中枢
先说结论:多个 Agent 一旦开始协作,它不是"多个工具排队执行",而是一张不断变化的临时组织网。没有中枢约束,这张网一定会乱。
1.1 三个绕不开的协作痛点
我最早接手过一个内容自动化项目,里面有策划 Agent、素材搜集 Agent、写作 Agent、审核 Agent,听起来分工很合理。真跑起来才发现,问题根本不在 Agent 本身的输出质量,而在协作链路上。
第一个痛点是消息乱序。Agent 之间是异步通信,A 发出了任务,B 处理完了回消息,但中间又插进来一条别的通知。某个 Agent 收到的不是"我预期的那条结果",而是"另一件事的状态变更",它的内部逻辑一旦没有鉴别能力,就会出现用旧数据写新文章这种低级错误。
第二个痛点是任务依赖。写作 Agent 需要等待素材搜集 Agent 返回三篇参考资料,策划 Agent 的结论没出来之前,素材搜集 Agent 不知道该收敛在哪几个话题维度上。这种依赖关系如果被硬编码在业务代码里,会变成一团 if-else 加上回调嵌套,根本维护不动。
第三个痛点是状态漂移。多个 Agent 各自维护自己的内存状态,没有统一的状态视图。任务跑到一半,某个 Agent 崩了、重启了,它的"第 3 步已完成"这种信息只存在于它的内存里。等你发现整个编排流程卡住了,要恢复现场就只能靠猜。
1.2 编排中枢到底在管什么
所谓编排中枢,本质上就是把上面三个痛点分别对应到三个确定性的技术模型上。
状态机负责回答"每个 Agent 现在处于什么阶段、能不能接受某个动作"。它把 Agent 的行为收敛到一个有限的、可预判的轨道上。DAG 负责回答"整个任务中,谁必须等谁,谁可以并行,谁的成败会影响下游谁"。它把业务依赖变成一张可检查、可调度、可重算的图。事件总线负责回答"Agent 之间通过什么形式通信、消息怎么路由、怎么保证不丢不重"。
这三者不是重叠的,而是各管一段。状态机管的是单点行为约束,DAG 管的是全局拓扑依赖,事件总线管的是连通性和消息可靠性。把这三样落地,多 Agent 系统的稳定性才会从"靠运气"变成"靠结构"。
2. 状态机:把 Agent 的行为约束在轨道上
状态机听起来学术味很重,但它不过就是一个非常朴素的道理:你想控制一个系统的复杂度,就要限制它的状态总数,并且明确什么条件能让它变到另一个状态。
2.1 状态机五要素与 Agent 生命周期
一个典型的状态机有五样东西:状态集合、事件集合、转移规则、动作、初始状态。放在 Agent 生命周期管理里,状态集合可以设计成这样:
- IDLE:Agent 已注册,空闲待命。
- READY:已领到任务,前置条件满足,可以开工。
- RUNNING:正在执行中。
- BLOCKED:执行过程中依赖的外部条件不满足,挂起等待。
- COMPLETED:任务正常完成。
- FAILED:任务执行出错、超时或业务校验不通过。
事件就是"触发状态变化的东西":比如收到任务、任务依赖就绪、执行超时、收到结果、收到重试指令。每个事件会让 Agent 从当前状态跳到一个新状态,非法的事件则直接拦截。比如 RUNNING 状态下又来一个 START 事件,应该抛异常,而不是默默忽略。
这个设计最关键的一点,是把 Agent 从"啥事都能干"变成"只有处于某状态才能干对应的事"。RUNNING 的 Agent 不能重复领取新任务,BLOCKED 的 Agent 不会硬着头皮继续跑。行为边界清晰了,异常排查时就少了一大半"这怎么跑到这里来了"的困惑。
2.2 状态转移表与落地代码
实际编码时我不会用复杂的框架,一个基于字典的状态转移表就够了。伪代码类似这样:
from dataclasses import dataclass from enum import Enum, auto class AgentState(Enum): IDLE = auto() READY = auto() RUNNING = auto() BLOCKED = auto() COMPLETED = auto() FAILED = auto() class AgentEvent(Enum): TASK_ASSIGNED = auto() DEPENDENCIES_READY = auto() EXEC_STARTED = auto() EXEC_FINISHED = auto() EXEC_TIMEOUT = auto() RETRY_REQUESTED = auto() TRANSITION_TABLE = { AgentState.IDLE: { AgentEvent.TASK_ASSIGNED: AgentState.READY }, AgentState.READY: { AgentEvent.DEPENDENCIES_READY: AgentState.RUNNING, AgentEvent.TASK_ASSIGNED: AgentState.READY # 幂等处理 }, AgentState.RUNNING: { AgentEvent.EXEC_FINISHED: AgentState.COMPLETED, AgentEvent.EXEC_TIMEOUT: AgentState.FAILED }, AgentState.BLOCKED: { AgentEvent.DEPENDENCIES_READY: AgentState.RUNNING, AgentEvent.RETRY_REQUESTED: AgentState.READY }, } class AgentStateMachine: def __init__(self, agent_id: str): self.agent_id = agent_id self.state = AgentState.IDLE def transition(self, event: AgentEvent): current = self.state allowed = TRANSITION_TABLE.get(current, {}) next_state = allowed.get(event) if next_state is None: raise IllegalTransitionError( f"Agent {self.agent_id} 当前状态 {current.name} 不接受事件 {event.name}" ) print(f"[状态机] Agent {self.agent_id}: {current.name} -> {next_state.name}") self.state = next_state这套写法的好处是转移规则集中、可视化简单、方便加审计日志。我实际项目里会在 transition 之前先做一次事件合法性校验,非法转移直接抛异常,而不是把错误数据带进下一步。
2.3 状态机设计的两个实操坑
第一个坑是状态设计得过细。有人把"等待素材1返回""等待素材2返回""等待素材3返回"做成三个独立状态,这让状态机完全失去了抽象价值。正确做法是统一收敛为 BLOCKED,然后依赖是否满足这件事交给 DAG 调度器去算,Agent 自己只需要知道"我现在卡住了"这个事实。
第二个坑是忽略超时状态。Agent 调 LLM 接口、外部搜索接口,都可能长时间无响应。状态机里如果没有超时事件,一个 RUNNING 状态可能会挂一整夜。我实际的策略是每个 RUNNING 都带一个最大执行时长,超时就走 EXEC_TIMEOUT 事件跳到 FAILED,然后由编排中枢决定是重试、降级还是终止。
3. DAG:让任务依赖关系一目了然
状态机管好单点行为之后,面临的下一个问题是:Agent 与 Agent 之间的任务依赖怎么组织。这条链路如果还是靠代码写死,项目超过十个节点之后一定会失控。DAG 才是这个场景的正解。
3.1 为什么是 DAG 而不是"串行链表"或"全并行"
有人会问:串行执行最简单,一个接一个跑不就行了?全并行不是更快吗?两者在复杂任务面前都不成立。
串行链表的致命问题是无法表达"并行分支"。素材搜集 Agent 完全可以拆成三个实例去搜集不同渠道的内容,它们之间没有依赖,完全可以同时跑。串行会把天然的并行度抹掉,拖慢整个任务。全并行的问题更明显:策划结论没出来,写作 Agent 不可能动笔;素材没攒够,审核无从谈起。没有依赖约束的全并行,就是无序竞争。
DAG 的价值在于它精确表达"谁必须先完成、谁可以在它完成后同时开工"。有依赖关系的任务形成边,没有依赖关系的任务形成并行分支,整张图既保留了并行度,又明确了先后顺序。
3.2 构建 DAG 的核心数据结构
工程上我习惯用邻接表来表达 DAG。一个任务节点就是一个 Agent 执行单元,一条边就是从上游节点指向下游节点的依赖关系。
from collections import deque class DAG: def __init__(self): self.nodes = {} # node_id -> NodeMeta self.out_edges = {} # node_id -> list[下游 node_id] self.in_edges = {} # node_id -> list[上游 node_id] def add_node(self, node_id: str, meta: dict): self.nodes[node_id] = meta self.out_edges.setdefault(node_id, []) self.in_edges.setdefault(node_id, []) def add_edge(self, upstream: str, downstream: str): self.out_edges.setdefault(upstream, []).append(downstream) self.in_edges.setdefault(downstream, []).append(upstream) def get_indegrees(self) -> dict: """返回每个节点当前入度,用于拓扑调度""" return {nid: len(self.in_edges[nid]) for nid in self.nodes} def is_acyclic(self) -> bool: """Kahn 拓扑排序,同时探测环路""" indeg = self.get_indegrees() queue = deque([nid for nid, deg in indeg.items() if deg == 0]) visited = 0 while queue: nid = queue.popleft() visited += 1 for downstream in self.out_edges[nid]: indeg[downstream] -= 1 if indeg[downstream] == 0: queue.append(downstream) return visited == len(self.nodes)DAG 的调度核心是拓扑排序,但是注意,工程里我不会等整张图排序完再执行,而是用"入度归零驱动"的方式:初始状态,所有入度为 0 的节点可以进入就绪队列;一个节点执行完成,就把它下游所有节点的入度减一;某个节点入度变成 0,说明它的全部前置依赖都完成了,它就可以被调度执行。
这种"动态就绪流"的好处不用说,它天然支持并行:就绪队列里的多个节点可以投递给不同的 Agent 实例同时执行。
3.3 动态 DAG 与运行时修改
最考验编排中枢的是"运行时改图"。真实场景里经常出现这种需求:策划 Agent 返回结果后,动态拆成 5 个素材搜集子任务,跑完后还要动态追加一个"素材质量汇总"节点。这要求 DAG 不只是建完就不可变,而是要支持增量式添加节点和边。
我踩过的坑是:新节点加入时没有检查是否引入了环路。比如素材汇总节点挂在策划节点下面,但策划节点在后续流程中又依赖素材汇总结果,这就形成了环,调度会卡死。解决办法不是每次加边都全图重算,而是只做增量环检测:从新边的下游节点出发,BFS 看能不能回到上游节点。能回到就拒绝加边。这个检查开销很小,但能避免整个 DAG 调度器卡在一个不可能的等待里。
4. 事件总线:Agent 之间不说悄悄话,全靠广播
状态机管了 Agent 自身的轨道,DAG 管了任务的依赖,第三个问题回到通信本身:Agent 之间的消息到底怎么传。我之前见过不少项目直接用函数调用让 Agent A 直接握住 Agent B 的引用,结果就是耦合成了蜘蛛网。事件总线就是为拆开这张网而生的。
4.1 事件驱动与请求响应:两种范式怎么选
Agent 之间通信有两条路线:请求响应式和事件驱动式。请求响应式看着自然,A 直接调用 B 的方法拿结果,但对异步、跨进程、多副本场景非常不友好。A 调 B 的时候 B 是不是活着?B 卡住了 A 怎么办?B 返回的结果迟到半天还该不该采纳?这些问题在请求响应模型里会变成无穷无尽的超时处理和重试策略。
事件驱动式是完全不同的思路。Agent 不直接"CALL 谁",只负责发布事件和订阅事件。就说素材搜集:它不是被写作 Agent 命令去干活,而是收到一个"素材搜集任务已创建"的事件,发现自己订阅了这类事件,就抢下来执行。执行完再发一个"素材搜集已完成"的事件。写作 Agent 关心的不是"谁去搜集了素材",而是"自己订阅的那类事件什么时候发生"。
两种方式可以混用,但我个人建议以事件驱动为主干:跨 Agent 的协作一律走事件,状态查询这类需要立即回执的动作再走轻量级请求响应。这样的系统,每个 Agent 不需要知道别人存在,一样能把活干成。
4.2 事件总线的核心接口设计
事件总线的核心抽象有三个东西:事件类型、发布者、订阅者。事件是事实记录,比如"任务启动""任务完成""任务失败""依赖已满足",它是已经发生的事情,不是命令。
from abc import ABC, abstractmethod from dataclasses import dataclass, field from typing import Any, Callable import asyncio @dataclass class Event: type: str source: str payload: dict = field(default_factory=dict) trace_id: str = "" seq: int = 0 class EventBus(ABC): @abstractmethod async def publish(self, event: Event): ... @abstractmethod def subscribe(self, event_type: str, handler: Callable): ... @abstractmethod async def start(self): ... @abstractmethod async def stop(self): ...一个简单但可用的内存事件总线实现,底层是 type -> handler 列表 的映射。发布事件时,事件总线把事件分发给所有订阅者。分布式的场景可以换成 Redis Stream 或 Kafka,把 EventBus 的实现替换掉即可,上层 Agent 代码不需要改动。这就是"面向接口编程"在工程上的实际回报。
4.3 事件总线可靠性的四条铁律
第一,消息不能丢。进程崩溃瞬间发出的事件,内存队列会直接丢。所以真实项目里事件总线要接一个持久化消息队列,事件先落盘再分发。
第二,消息尽可能不重。分布式系统想做到绝对不重复分发,代价非常高。实际落地方案是用消息 ID 加幂等消费:每个事件带一个全局唯一的 trace_id 加 seq,消费者端保存已处理过的 ID 集合,重复消息直接忽略。
第三,顺序要有边界。全局严格有序在多 Agent 并发的场景下是伪需求。我只要保证同一个任务派生的所有事件有序即可:事件里带 seq,消费者有一个小的重排窗口。跨任务的全局顺序反而不重要。
第四,消费端要背压。事件一旦爆发式推送,下游 Agent 处理不过来,内存就会涨。我一般会给每个订阅者的待处理队列设置上限,满了就暂停拉取,形成天然的背压,而不是把消息无限堆在内存里。
5. 编排中枢完整实现:状态机、DAG、事件总线如何协同
三个组件各有分工,但它们不是三个孤立模块,而是通过明确的数据流接在一起。这一节我用一条完整任务链路来演示它们怎么配合。
5.1 一次任务的全生命周期走读
场景还是内容项目。任务从用户提交选题开始。
第一步,创建一个 DAG 实例,初始只有两个节点:策划节点和素材搜集前置节点。事件总线发布一个"任务已启动"事件。策划 Agent 订阅了这个事件,收到后开始执行。
策划 Agent 内部有一个状态机,此刻从 IDLE 跳到 RUNNING。它调用 LLM 产出内容策略。完成后,它的状态机跳到 COMPLETED,同时通过事件总线发布"策划已完成"事件。
这个事件是 DAG 调度器订阅的关键信号。调度器收到后,发现自己维护的 DAG 里,策划节点已经完成,于是把策划节点的所有下游节点入度减一。如果此时候素材节点入度归零,就进入就绪队列。调度器从就绪队列取出素材任务,再发布一个"素材搜集任务已创建"事件。
三个素材 Agent 实例都订阅了这类型事件,各自领走一个子任务,并行搜集。它们每个都维护自己的状态机,完成时各发各的"素材子任务已完成"事件。调度器统计到三个子任务全部完成后,汇总节点入度归零,发布"素材搜集全部完成"事件。
写作 Agent 收到事件后,开始动笔。写完发布"初稿已完成"事件。审核 Agent 收到事件进入审核。如果审核不通过,发布"审核未通过"事件,写作 Agent 收到后状态机从 COMPLETED 回到 READY,DAG 调度器把这个追加的新一轮写作节点加入 DAG,开始下一轮循环。整个过程,没有一个 Agent 直接调用另一个 Agent,全靠事件总线传递事实;每个 Agent 每个时刻处于哪个状态,状态机给了明确答案;任务下一步该跑谁,DAG 给了准确依据。
5.2 关键数据结构:把三个组件粘在一起
为了让三个组件协同,我通常还会加一个 Orchestrator 对象,它自己持有一张"任务实例 -> DAG + Agent状态机 + 事件处理器"的映射。
class Orchestrator: def __init__(self, bus: EventBus): self.bus = bus self.tasks = {} # task_id -> TaskRuntime async def handle_task_created(self, event: Event): task_id = event.payload["task_id"] dag = build_initial_dag(task_id) runtime = TaskRuntime(task_id=task_id, dag=dag) self.tasks[task_id] = runtime await self.bus.publish(Event(type="task.started", source="orchestrator", payload={"task_id": task_id}, trace_id=task_id)) async def handle_agent_completed(self, event: Event): task_id = event.trace_id runtime = self.tasks[task_id] node_id = event.payload["node_id"] runtime.mark_node_completed(node_id) ready_nodes = runtime.pop_ready_nodes() for node in ready_nodes: await self.bus.publish(Event(type="task.node_ready", source="orchestrator", payload={"node_id": node}, trace_id=task_id))这个例子展示了事件总线的核心衔接作用:每个 Agent 完成事件进来,Orchestrator 只要更新 DAG 状态、取出新就绪节点,再发布新事件即可。它的状态机管理被封装在 Agent 侧,DAG 调度在 Orchestrator 侧,二者通过事件总线完成握手。
5.3 一个避坑指南:谁来触发 Agent 工作
新手最容易搞混的问题是谁负责"叫醒"Agent。我在实战中的结论是:Agent 不应该被轮询唤醒,也不应该被显式点名,而是由事件总线按订阅分发。但是,订阅模式要防止一个事件被多个同类 Agent 重复消费后导致重复执行。
解决方法是引入"work-queue 型事件":事件带 work_id,多个候选 Agent 收到后抢锁,只有一个能抢到执行权。抢不到的直接忽略。这个细节如果不处理,素材 Agent 配置了三副本,一个"素材子任务已创建"事件发出来,三个副本同时开跑,就会产生三份重复的搜集结果。
6. 常见问题与排查技巧实录
这块内容全部来自真实排查过的问题列表,比理论更实用。我按问题、原因、排查手段和解决建议整理成速查式写法,踩过坑的同学可以直接对着查。
6.1 编排任务卡死:某个 Agent 永远在 RUNNING
现象是 DAG 里有几个下游节点迟迟不进入就绪队列,日志显示某个上游 Agent 状态一直是 RUNNING。
排查第一步先看这个 Agent 最近有没有发过"完成"事件。大概率是没有。接着看它调用的外部接口是不是超时时间设得太大,或者 LLM 调用一直在重试。最隐蔽的原因是我之前提过的:状态机忘了设计超时转移,RUNNING 没有出口。
解决建议:给所有 RUNNING 状态统一加最大执行时长;外部接口超时时间至少要小于状态机超时时间,这样状态机先收到超时事件,才能主导失败逻辑,而不是等外部调用一直挂着。
6.2 DAG 环路导致调度死循环
现象是调度器日志里大量输出"节点已就绪、节点已执行",但某些节点永远执行不到;或者拓扑排序抛异常说检测到环。
环通常来自运行时动态加边。最典型的就是我前面提到的"重试依赖反向":审核失败后要回写作,如果代码直接加一条从审核指向写作的边,同时原有从写作指向审核的边还在,就形成了 写作->审核->写作 的环。
解决建议:动态添加边必须做增量环检测;重试流程不要在原 DAG 上改边,而是把原节点视为完成、生成一个新的"修订写作"节点挂到终点后面,再让它通向审核。这样既表达了重试,又不破坏 DAG 的无环性。
6.3 事件风暴:Agent 集体过载
现象是事件总线队列积压,Agent 端处理延迟暴涨,系统整体吞吐骤降。这在业务上常见的触发点是"大量任务同时被创建",比如运营批量导入几百个任务,每个任务拆成 8 个节点,瞬间产生几千个事件。
解决建议有两个层面。第一靠背压,订阅者队列有界,忙不过来就暂停拉取。第二靠批量合并,比如素材子任务完成事件不要一次性发几十个,而是聚合成一批"多事件完成"消息再发。窗口大小通常取几百毫秒,吞吐翻倍但延迟没明显影响。
6.4 幂等消费:重试导致的重复执行
现象是任务日志里发现同一个节点执行了两次,两次结果还会互相覆盖。事件总线重试机制天然会带来重复消息,如果消费者不做幂等,重复就是大概率事件。
解决建议:每个带业务含义的事件必须包含幂等键,消费者侧记录已处理 ID 集合。另外我还习惯在结果写入时做"版本号比对":消息里带上来源节点版本号,下游更新时校验目标版本,防止旧消息覆盖新消息。这个双保险在我项目里救过不止一次。
排查这块再补一个我自己的小技巧:给事件总线加一个"黑暗模式"开关,开启后所有事件只记录不投递。这能快速判断一个 Agent 是"没触发"还是"触发了但没执行完",是我用它排查诡异问题的首选手段。
7. 收尾的一点个人经验
这套状态机加 DAG 加事件总线的骨架,我迭代了好几版才稳定下来。回头复盘,最早犯的错是把三个组件揉在一起做了一个"万能中枢",结果状态机的状态里塞了 DAG 依赖信息、事件总线里裹了状态迁移逻辑,最后谁都不好改。后来拆清楚"状态机管单点行为、DAG 管全局依赖、事件总线管通信"之后,系统才真正变得干净。
如果让我给正在做类似系统的同学一个建议,我会说:不要一开始就追求最复杂的分布式事件总线,先在单机内把这三者的数据流跑通,再用真实任务打磨状态转移表和 DAG 结构。等稳定了,再逐步把事件总线换成 Redis Stream 或 Kafka 这类持久化消息组件,整个过程会顺畅得多。多 Agent 项目能走多远,不取决于模型有多聪明,而取决于你给它铺的这套"高速公路"有多稳。