在这里插入图片描述

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 被阻塞”。不是模糊的"出错了",而是清晰的因果链。

三个算法一句话回顾

  1. DFS 染色 —— 白灰黑三色标记,发现灰色节点即检测到环,输出人类可读的环路径
  2. Kahn 分层 —— BFS 逐层剥离入度为 0 的节点,每一层就是一个并行批次
  3. 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。

Logo

小龙虾开发者社区是 CSDN 旗下专注 OpenClaw 生态的官方阵地,聚焦技能开发、插件实践与部署教程,为开发者提供可直接落地的方案、工具与交流平台,助力高效构建与落地 AI 应用

更多推荐