14 个 AI 同时工作,我是怎么不让它们打架的

TL;DR(30 秒速览)
- 14 个 AI 分析师组成投资委员会,有的看技术面、有的看基本面、还要辩论——怎么编排?答案是 DAG(有向无环图)
- 三个核心算法:DFS 白灰黑染色检测循环依赖、Kahn 算法分层拓扑排序、BFS 级联标记失败下游
Semaphore有界并发 +AtomicBoolean暂停/恢复 +Mutex保护状态——全 Kotlin 协程实现- 三种任务类型(SINGLE / TEAM / DELIBERATION)在 DAG 调度层完全透明——关注点分离
- 开源地址:github.com/haibingzhao/easyai
前情提要:在上一篇文章中,我们聊了为什么手写 ReAct 循环——绕过 Spring AI 的
ToolCallingAdvisor,直接控制 Agent 循环的每一步。但单个 Agent 搞定了,多个 Agent 怎么协作?
问题比想象的复杂
想象这个场景:14 个 AI 分析师组成一个投资委员会,执行一次完整的投资分析:
数据采集 → 4 路并行分析 → 多空辩论 → 交易计划 → 风险辩论 → 最终决策
这不是"Agent A 交给 Agent B"的简单链式调用。它是一个有依赖关系的 DAG(有向无环图):
- 4 个分析师可以同时工作(并行)
- 辩论必须等所有分析完成(汇聚)
- 交易计划依赖辩论结果(顺序)
- 一个分析师失败了,依赖它的下游全部要被阻塞
核心矛盾: 多 Agent 编排不是"谁来调用谁"的问题,而是"哪些可以并行、哪些必须等待、失败了怎么办"的图论问题。
三个算法搞定 DAG 编排

一句话概括: DAG 编排的本质是三个问题——有没有环?执行顺序是什么?失败了谁受影响?分别对应三个图算法。
算法 1:DFS 环检测(白-灰-黑染色法)
用户在可视化编辑器里拖拽连线时,可能不小心创建循环依赖:A 等 B、B 等 C、C 又等 A。必须检测并拒绝。
白-灰-黑染色法是经典的 DFS 环检测算法:
| 颜色 | 含义 |
|---|---|
| 白色 | 未访问 |
| 灰色 | 正在访问中(在递归栈里) |
| 黑色 | 已完全处理 |
关键逻辑:如果在 DFS 过程中遇到了一个灰色节点,说明找到了环——因为灰色节点意味着"我还在处理中,你怎么又回来了?"
// DagAlgorithms.kt — 简化版
fun dfs(nodeId: String, path: MutableList<String>) {
color[nodeId] = Color.GRAY // 标记为"正在处理"
path.add(nodeId)
for (neighbor in adjacency[nodeId].orEmpty()) {
when (color[neighbor]) {
Color.GRAY -> {
// 遇到了灰色节点 → 发现环!
val cycleStart = path.indexOf(neighbor)
val cycle = path.subList(cycleStart, path.size) + neighbor
throw IllegalStateException(
"DAG cycle detected: ${cycle.joinToString(" → ")}"
)
}
Color.WHITE -> dfs(neighbor, path) // 继续深入
Color.BLACK -> { /* 已处理,跳过 */ }
}
}
path.removeLast()
color[nodeId] = Color.BLACK // 标记为"处理完毕"
}
亮点:人类可读的环路径报告。 不是简单抛出"存在环",而是输出 A → B → C → A——用户一眼就知道哪里连错了。
算法 2:Kahn 算法分层拓扑排序
检测完环,下一个问题是:哪些任务可以并行?
普通的拓扑排序只给出一个线性顺序(A → B → C),但我们需要的是分层——同一层的任务互不依赖,可以并行执行。这就是 Kahn 算法的变体:
初始化:计算每个节点的入度(有多少个上游依赖)
循环:
1. 找出所有入度为 0 的节点 → 它们组成当前层
2. 移除这些节点,更新下游的入度
3. 重复,直到所有节点处理完毕
对应到投资分析的例子:
Layer 0: [AnalystStater] ← 1 个任务(数据采集)
Layer 1: [Market, News, Sentiment, ← 4 个任务(并行分析)
Fundamentals]
Layer 2: [BullBearDebate] ← 1 个任务(辩论,依赖 Layer 1 全部完成)
Layer 3: [Trader] ← 1 个任务(交易计划)
Layer 4: [RiskDebate] ← 1 个任务(风险辩论)
Layer 1 的 4 个分析师互不依赖,同时启动。Layer 2 的辩论必须等 Layer 1 全部完成。
// DagAlgorithms.topologicalLayers() — 简化版
fun topologicalLayers(tasks: List<SwarmTask>): List<List<String>> {
val inDegree = mutableMapOf<String, Int>()
tasks.forEach { inDegree[it.id] = it.dependsOn.size }
val layers = mutableListOf<List<String>>()
while (processed < tasks.size) {
// 入度为 0 → 没有待完成的依赖 → 可以执行
val layer = inDegree.filter { it.value == 0 }.keys.toList()
if (layer.isEmpty()) {
throw IllegalStateException("DAG cycle detected: ...")
}
layers.add(layer)
// 移除已处理节点,更新下游入度
for (nodeId in layer) {
inDegree.remove(nodeId)
for (dependent in adjacency[nodeId].orEmpty()) {
inDegree[dependent] = inDegree[dependent]!! - 1
}
}
}
return layers
}
算法 3:BFS 失败级联
一个上游任务失败了,所有依赖它的下游任务不应该继续等待——它们应该被精确标记为 BLOCKED。
不是简单取消,而是 BFS 遍历:从失败节点出发,沿着邻接表逐层标记所有传递依赖。
// DagAlgorithms.resolveDependencies() — 失败级联
if (failed) {
val toBlock = ArrayDeque<String>()
adjacency[completedTaskId]?.forEach { toBlock.add(it) }
while (toBlock.isNotEmpty()) {
val taskId = toBlock.removeFirst()
val task = tasks.find { it.id == taskId } ?: continue
if (task.status == SwarmTaskStatus.PENDING) {
task.status = SwarmTaskStatus.BLOCKED
task.blockedBy = listOf(completedTaskId)
// 继续向下游传播
adjacency[taskId]?.forEach { toBlock.add(it) }
}
}
}
亮点: 用户在前端可以看到完整的级联路径——“Market Analyst 失败了 → BullBearDebate 被阻塞 → Trader 被阻塞 → RiskDebate 被阻塞”。不是模糊的"出错了",而是清晰的因果链。
三个算法一句话回顾
- DFS 染色 —— 白灰黑三色标记,发现灰色节点即检测到环,输出人类可读的环路径
- Kahn 分层 —— BFS 逐层剥离入度为 0 的节点,每一层就是一个并行批次
- BFS 级联 —— 从失败节点 BFS 标记所有下游为 BLOCKED,精确而非暴力取消
并发控制:为什么是 Semaphore 而不是 Actor

DAG 调度层确定了"哪些任务可以并行",但 14 个 Agent 同时调用 LLM API 会触发限流。需要有界并发。
EasyAI 的选择是 Semaphore(maxConcurrency):
class SwarmRuntime(...) {
private val semaphore = Semaphore(maxConcurrency) // 默认 4
private val dagMutex = Mutex()
// 每个任务执行时获取信号量
semaphore.withPermit {
// 执行 Agent...
}
}
| 并发原语 | 用途 | 为什么选它 |
|---|---|---|
Semaphore(4) |
限制同时运行的 Agent 数量 | 有界并发,防止 API 限流 |
AtomicBoolean |
abort / pause 信号 | 无锁读写,协程友好 |
Mutex |
保护 DAG 状态的并发修改 | 任务完成时更新入度表需要互斥 |
ConcurrentHashMap |
存储任务摘要 | 层间传递,读写都线程安全 |
为什么不用 Kotlin Actor? Actor 模型有自己的线程调度策略,和 Spring 的虚拟线程模型(spring.threads.virtual.enabled=true)不一定兼容。Semaphore + Mutex 是最朴素的选择,也是最可预测的。
暂停和恢复:不只是取消
生产环境中,用户可能需要:
- 暂停:先看看 Layer 1 的分析结果,再决定要不要继续
- 恢复:看完结果后从中断点继续执行
- 取消:彻底终止
for ((layerIndex, layer) in layers.withIndex()) {
// 每层开始前检查
if (abortSignal()) {
markRemainingTasksCancelled(run)
run.status = SwarmRunStatus.CANCELLED
break
}
if (pauseRequested?.get() == true && layerIndex > 0) {
run.status = SwarmRunStatus.PAUSED
break // 当前层执行完后停止
}
// 执行本层所有任务...
}
恢复时,从数据库加载已完成任务的状态,找到第一个未完成的层,从那里继续:
// SwarmRuntime.resume() — 找到第一个未完成的层
val layers = DagAlgorithms.topologicalLayers(run.tasks)
var resumeFromLayer = 0
for ((layerIndex, layer) in layers.withIndex()) {
val allDone = layer.all { taskId ->
val task = run.tasks.find { it.id == taskId }
task != null && (task.status == COMPLETED || FAILED || BLOCKED)
}
if (!allDone) { resumeFromLayer = layerIndex; break }
}
// 从 resumeFromLayer 继续执行
任务类型路由:调度层不关心你怎么干活
DAG 调度层只负责"什么时候执行谁",不关心"怎么执行"。三种任务类型各有自己的执行器:
SwarmRuntime(DAG 调度层)
├── SINGLE → SwarmWorkerExecutor(单 Agent 执行)
├── TEAM → TeamTaskExecutor(Leader-Member 协作)
└── DELIBERATION → DeliberationTaskExecutor(多轮辩论 + 裁判)
// SwarmRuntime.executeTask() — 按类型路由
val result = when (task.type) {
TaskType.SINGLE -> workerExecutor.runSingleWorker(task, run, ...)
TaskType.DELIBERATION -> deliberationExecutor.runDeliberation(task, run, ...)
TaskType.TEAM -> teamExecutor.runTeam(task, run, ...)
}
每种任务类型的产出物不同(单 Agent 返回文本、辩论返回多轮记录、TEAM 返回成员执行记录),但 DAG 调度层只关心一件事:你完成了没有?摘要是什么?
这就是关注点分离——调度层不需要知道辩论是怎么进行的,它只需要知道辩论的结果,然后传递给下游。
上游摘要传递:下游不重复做上游的工作
下游任务的 Prompt 通过 Jinja2 模板引用上游任务的执行摘要:
taskSummaries: ConcurrentHashMap<String, String>
任务完成 → 摘要存入 taskSummaries → 下游任务的 promptTemplate 通过 {{ task_id }} 引用
例如,交易计划任务的 promptTemplate:
上游辩论结果:
{{ BullBearDebate }}
请基于以上辩论结果,制定结构化交易计划。
{{ BullBearDebate }} 在运行时被替换为辩论任务的执行摘要。下游不需要重新分析市场数据——上游已经做过了。
inputFrom 机制还支持别名映射——下游可以用语义化的名字引用上游:
{
"id": "trade_plan",
"inputFrom": { "debate_result": "BullBearDebate" },
"promptTemplate": "辩论结果:{{ debate_result }}\n请制定交易计划。"
}
实战案例
案例 1:14 Agent 投资分析
Layer 0: [AnalystStater] ← 数据采集(1 个 Agent)
Layer 1: [Market] [News] [Sentiment] [Fund.] ← 4 路并行分析(4 个 Agent)
Layer 2: [BullBearDebate] ← 多空辩论(3 个 Agent,DELIBERATION)
Layer 3: [Trader] ← 交易计划(1 个 Agent)
Layer 4: [RiskDebate] ← 风险辩论(4 个 Agent,DELIBERATION)
5 层、9 个任务、14 个 Agent。Layer 1 的 4 个分析师并行执行,Semaphore 控制同时只有 4 个 Agent 调用 LLM API。总耗时 10-20 分钟。
案例 2:6 Member 编程团队
一个 TEAM 类型的 Agent,Leader 自动分解任务并协调 6 个专业成员:
| 成员 | 职责 | 工具权限 |
|---|---|---|
| Researcher | 代码定位、依赖分析 | read, grep, bash, websearch |
| Full-Stack Engineer | 代码实现 | read, write, edit, bash |
| QA | 测试验证 | read, bash |
| Code Reviewer | 代码审查(只读) | read, grep |
| UI Operator | 浏览器端到端验证 | read, bash + browser-use MCP |
| Debug Engineer | 故障诊断 | read, bash, grep |
Leader 先派 Researcher 定位问题范围,再派 Engineer 实现,QA 和 Reviewer 并行验证——这本身也是一个微型的 DAG。
踩坑记录
三个坑,每个都花了至少一天调试。
坑 1:Kahn 算法的入度表必须动态更新。 最初的实现是在执行前一次性计算所有层,然后按层执行。但任务完成后需要更新下游的 blockedBy 列表——如果入度表是静态的,下游任务不知道上游已经完成。修复方案:每层执行完后,通过 resolveDependencies() 动态更新。
坑 2:失败任务要及时释放 Semaphore。 一个任务超时失败后,如果它的下游没有被及时标记为 BLOCKED,它们会一直等待 Semaphore 的 permit——而 permit 已经被占满了。修复方案:任务失败时立即调用 resolveDependencies(failed = true) 级联标记下游,释放资源。
坑 3:并发持久化必须立即做。 任务完成后如果不立即写入数据库,进程崩溃时会重复执行已完成的任务。修复方案:每个任务完成后立即 store?.saveTask(run.id, task),在 NonCancellable 块中确保持久化不被取消中断。
写在最后
多 Agent 编排的核心不是"谁来调用谁",而是图论。
DFS 染色保证无环、Kahn 算法确定并行层级、BFS 级联处理失败传播——三个经典算法组合起来,就是一个生产级的 DAG 编排引擎。
| 维度 | 简单链式调用 | DAG 编排 |
|---|---|---|
| 并行能力 | 无(严格顺序) | 同层并行、层间顺序 |
| 依赖表达 | 只能 A → B | 任意 DAG 拓扑 |
| 失败处理 | 整体重试 | 精确级联标记 BLOCKED |
| 暂停/恢复 | 不支持 | 层间暂停、断点恢复 |
| 任务类型 | 单一 | SINGLE / TEAM / DELIBERATION 混合 |
EasyAI 的 Swarm 模块用不到 600 行 Kotlin 代码实现了完整的 DAG 编排——核心算法 DagAlgorithms 只有 179 行。好的算法不需要很多代码,但需要想清楚。
下一篇,我们聊聊 Deliberation——如何让 AI 像法庭一样对抗辩论,裁判自主判断何时收敛,输出更接近真相的决策。
开源地址:https://github.com/haibingzhao/easyai
欢迎 Star、Issue 和 PR。
更多推荐



所有评论(0)