open-multi-agent 调度原理完整拆解:事件驱动执行器、TaskQueue 依赖图与 AgentPool 并发控制
【免费下载链接】open-multi-agentTypeScript AI agent orchestration framework with dynamic workflows. Describe the goal, not the graph: a coordinator plans the task DAG at runtime and runs it on any LLM (Claude, ChatGPT, Gemini, DeepSeek, or local models).项目地址: https://gitcode.com/gh_mirrors/op/open-multi-agent
open-multi-agent是一个用 TypeScript 编写的 AI Agent 编排框架(多智能体编排框架):你只描述目标,协调器(Coordinator)会在运行时把任务规划成一张有向无环图(DAG),交给 Claude、ChatGPT、Gemini、DeepSeek 或本地模型等任意 LLM 并行执行。🧭 本文将深入拆解它内部的三块调度核心:事件驱动执行器、带依赖图的TaskQueue,以及负责并发控制的AgentPool,帮你彻底看懂一个任务从"就绪"到"完成"的全过程。
30 秒看懂整体架构:谁负责什么?
在 open-multi-agent 中,一次团队执行(runTeam())的调度链路可以概括为三层,各司其职:
| 组件 | 源码位置 | 职责 |
|---|---|---|
| TaskQueue(任务队列) | packages/core/src/task/queue.ts | 持有所有任务,维护依赖图与状态机,发出调度事件 |
| 事件驱动执行器 | packages/core/src/orchestrator/task-execution.ts | 订阅队列事件,决定"现在该派哪个任务" |
| AgentPool(Agent 池) | packages/core/src/agent/pool.ts | 用信号量给全体 Agent 限流,保证并发不超载 |
| Scheduler(调度器) | packages/core/src/orchestrator/scheduler.ts | 把就绪任务分配给最合适的 Agent |
这个分层设计的核心思想是:队列管"能不能跑",执行器管"该不该派",池子管"同时能跑几个"。三者解耦后,任何一层都可以独立演进——官方文档 docs/task-scheduling.md 中也明确写道:AgentPool的信号量始终是并发的最终权威(concurrency authority)。
事件驱动执行器:任务就绪即调度,告别轮询
很多工作流框架靠"轮询检查进度"来推进流程,而 open-multi-agent 默认采用事件驱动模式(Event-driven execution)。执行器executeQueue()在启动时做三件事(见 task-execution.ts):
- 订阅
TaskQueue的task:ready事件——某任务的所有依赖一完成,队列立即发出该事件,执行器无需轮询; - 维护两个集合:就绪集合(readyTaskIds)和在飞映射(inFlight Map),记录正在执行的任务;
- 每轮循环通过**派发门(dispatch gate)**检查四件事:是否被调用方取消(abort)、是否超出 token 预算、是否超过 AgentPool 容量、是否有待审批的边界。✅
只有全部通过,任务才会被派发给 AgentPool。这种"事件喂料 + 门控放行"的循环还有一个巧妙细节:即使某个下游任务已经就绪,执行器也会等它的前置任务彻底落盘(结果、检查点、追踪事件写完)后再启动,保证"完成事件早于开始事件"的时序一致,方便你事后用观测面板复盘。🔍
失败与跳过如何"级联"?
依赖图的另一半逻辑是失败传播,全部在 queue.ts 中完成:
fail():某任务失败时,递归地把所有传递性依赖它的下游任务标记为 failed,避免它们永远卡在 blocked 状态;skip():审批被拒绝等场景下,skipRemaining()会先停止派发、等在飞任务排空,再把剩余任务统一跳过;- 关键在于:级联只影响下游分支,无关分支继续并行跑,不会"一损俱损"。
TaskQueue 依赖图:任务状态机与 5 种核心事件
TaskQueue是所有任务的"单一事实来源"(single source of truth)。每个任务在六种状态间流转,队列以事件方式对外广播变化:
| 状态 | 含义 | 触发时机 |
|---|---|---|
pending | 就绪,等待派发 | 无依赖,或所有依赖已完成 |
blocked | 被依赖阻塞 | 存在未完成的dependsOn任务 |
in_progress | 正在执行 | 执行器派发后 |
completed/failed/skipped | 三种终态 | 执行结束 / 失败级联 / 跳过 |
依赖解锁:unblockDependents()的 O(n) 扫描
任务完成时,队列会调用 unblockDependents():扫描所有 blocked 任务,凡依赖链全部满足者立即"晋升"为 pending,并为每个新解锁的任务发出task:ready事件——这就是事件驱动执行器被"叫醒"的信号。实现上,任务数组和 ID 索引 Map 各只构建一次,把整个扫描控制在 O(n) 而非 O(n²)。
任务队列发出的 5 种事件
队列的对外接口非常克制,只有 5 种命名事件:
task:ready— 新任务就绪(含依赖解锁)task:complete— 任务完成,随后依次触发其下游的task:readytask:failed/task:skipped— 失败/跳过,并携带级联信息all:complete— 全部任务到达终态,执行器可以收尾 🎉
订阅方式也很简单:queue.on('task:ready', handler)返回一个退订函数,幂等安全,详见 queue.ts 事件小节。
快照与恢复:断点续跑的基础
生产场景必须考虑崩溃恢复。TaskQueue 支持snapshot()全量序列化、fromSnapshot()精确重建,且可传入resetInProgress: true把"崩溃时正在执行"的任务重置为可重跑状态。配合 memory/checkpoint.ts 的检查点机制,一次被中断的runTeam()可以从中断处继续,而不是从头烧 token。💰
AgentPool 并发控制:信号量如何给 AI 团队"限流"
LLM 调用又贵又慢,并发失控意味着账单爆炸和 API 限流。open-multi-agent 用一把自研计数信号量(utils/semaphore.ts)给 AgentPool 限流,构造时默认maxConcurrency = 5。
AgentPool 的并发控制其实是两把锁,这一点初学者最容易忽略:
- 池级信号量(
Semaphore(maxConcurrency)):限制整个池子同时运行的 Agent 数量。超出上限的调用会在acquire()中排队,FIFO 依次放行; - Agent 级互斥锁(每个 Agent 一把
Semaphore(1)):同一个 Agent 实例内部的status、messages、tokenUsage 是可变状态,两件事同时打给它会互相踩脚,所以同一 Agent 的运行被串行化。⚔️
一个值得学习的工程细节是加锁顺序:run()中先拿 Agent 锁、再拿池级信号量(见 pool.ts),这样第二个打到同一 Agent 的调用会在 Agent 锁处等待,不占用池的并发槽位,避免"占着茅坑"式的资源浪费。
委托(Delegation)为什么走另一条路?
Agent 可以调用delegate_to_agent把子任务委托给同队伙伴。此时走的是runEphemeral():为被委托方新建一个一次性 Agent 实例,只拿池级信号量、跳过 Agent 级锁。源码注释里解释得很直白——如果不这样做,A 委托 B 时 B 又委托 A,双方各持对方的 Agent 锁,就会互相死锁。🕳️
此外,池还暴露了availableRunSlots属性:执行器在派发前会先检查剩余槽位(inFlightCount >= runConcurrencyLimit即暂停),确保"委托一个任务"永远不会把池子挤到死锁边缘。
三种运行入口怎么选?
| 方法 | 适用场景 | 并发约束 |
|---|---|---|
run(name, prompt) | 单个指定 Agent | Agent 锁 + 池信号量 |
runParallel(tasks) | 一批任务并行扇出 | 池信号量封顶,失败转为错误结果而非抛异常 |
runAny(prompt) | 不指定 Agent,轮询分派 | Agent 锁 + 池信号量 |
调度器 5 大策略:为不同 Agent 团队挑选"排班方式"
任务就绪后,"派给谁"由 Scheduler 决定。内置 5 种策略,dependency-first是默认:
round-robin— 按索引轮流分派,适合能力完全对等的 Agent;least-busy— 派给当前在跑任务最少的 Agent,适合任务耗时差异大;capability-match— 先按硬性要求过滤,再按能力/关键词亲和度打分,适合角色分工明确的团队;dependency-first(默认)— 优先执行"关键路径"上的任务,即解锁下游最多的任务,靠正向 BFS 统计每个任务的"关键度",特别适合依赖密集的工作流;composite— 加权组合关键度、能力匹配与当前负载(默认权重 fit 0.7 / load 0.3),多目标综合排序。
举个直觉例子:如果你的 DAG 是"调研 → 写初稿 → 三人并行评审 → 汇总",dependency-first会先把"调研"推出去,因为它卡着后面 4 个任务;评审三兄弟则自然并行。
总结:一张链路看懂 open-multi-agent 的调度
把本文三块内容串起来,一次任务派发就是下面这条流水线:
依赖完成→ TaskQueue 发
task:ready→ 执行器就绪集合更新 → 派发门检查(取消/预算/容量/审批)→ Scheduler 选 Agent → AgentPool 两把锁放行 → 执行、重试、验证 → 结果回写队列,task:complete再次唤醒循环 🔄
这正是 open-multi-agent 的设计哲学:"描述目标,而非画图"。依赖关系、并发上限、失败传播全部由框架在运行时托管,你只关心每个任务做什么、谁能做。想动手验证,可以从 examples/basics/team-collaboration.ts 这类最小示例入手,再对照 packages/core/examples/patterns/event-driven-dag.ts 看事件驱动 DAG 的完整用法。
【免费下载链接】open-multi-agentTypeScript AI agent orchestration framework with dynamic workflows. Describe the goal, not the graph: a coordinator plans the task DAG at runtime and runs it on any LLM (Claude, ChatGPT, Gemini, DeepSeek, or local models).项目地址: https://gitcode.com/gh_mirrors/op/open-multi-agent
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考