FastGPT 工作流 Runtime 引擎解析:基于 Tarjan SCC 与 DFS 边分类的节点调度与执行机制
2026/9/11 11:55:06 网站建设 项目流程

FastGPT 工作流 Runtime 引擎解析:基于 Tarjan SCC 与 DFS 边分类的节点调度与执行机制

【免费下载链接】FastGPTFastGPT is a knowledge-based platform built on the LLMs, offers a comprehensive suite of out-of-the-box capabilities such as data processing, RAG retrieval, and visual AI workflow orchestration, letting you easily develop and deploy complex question-answering systems without the need for extensive setup or configuration.项目地址: https://gitcode.com/GitHub_Trending/fa/FastGPT

FastGPT 的可视化 AI 工作流本质上是一个有向图执行引擎,用户在画布上编排的"节点 + 连线"会在运行时被编译为带状态的有向图,由WorkflowQueue统一调度。本文以仓库设计文档 runtime.md 为主体,结合 dispatch/index.ts 与 tarjan.ts 的源码实现,系统讲解 FastGPT 工作流 Runtime 如何通过 Tarjan 强连通分量算法识别循环、通过 DFS 边分类区分回边,并基于"节点边分组 + 边状态机"完成分支、循环、并行、工具调用等复杂拓扑的正确调度。读完本文,你将掌握 Runtime 的核心数据结构、边分组算法、节点运行状态判定规则,以及它在典型工作流拓扑下的行为与测试覆盖情况。

一、Runtime 整体架构

FastGPT 工作流 Runtime 的设计目标是:把画布上"任意可连"的有向图(允许分支、循环、并行、交叉)变成一组语义清晰、可判定的执行规则。整个执行体系围绕三块内容展开:

  1. 图论分析层:使用 Tarjan 算法找出强连通分量(SCC),判断节点是否处于循环;使用 DFS 对每条边做分类(树边 / 回边 / 前向边 / 跨边),识别循环边。
  2. 边分组层:为每个节点的所有输入边构建"边分组"(NodeEdgeGroups),把复杂的图结构归约为若干组"或/且"关系的判定条件。
  3. 队列调度层WorkflowQueue维护活跃节点队列(activeRunQueue)与跳过节点队列(skipNodeQueue),用迭代循环取代递归,配合并发上限持续驱动节点执行,直到所有路径收敛。

1.1 WorkflowQueue:核心执行类

WorkflowQueue定义于 dispatch/index.ts,是工作流执行的心脏,负责管理节点执行队列与状态。其关键属性包括:

属性作用
runtimeNodesMap节点 ID 到运行时节点对象的映射(初始化于 index.ts)
edgeIndex边的索引,按 source 和 target 两个维度预构建(bySource/byTarget
nodeEdgeGroupsMap预构建的"节点 → 输入边分组" Map,运行时直接查询,避免重复计算
activeRunQueue活跃运行队列,存放待检查/待运行的节点 ID
skipNodeQueue跳过节点队列,存放整条被跳过路径上的节点

关键方法一览:

  • buildEdgeIndex():扫描全部运行时边,构建bySourcebyTarget两个索引 Map(实现于 index.ts)。
  • buildNodeEdgeGroupsMap():对每个节点执行"DFS 边分类 → Tarjan SCC 分析 → 按分支句柄分组"三步流程,一次性构建全图边分组(index.ts)。
  • getNodeRunStatus():根据预构建的边分组判定节点应runskip还是wait(index.ts)。
  • addActiveNode():把节点加入活跃队列,若当前没有处理循环则触发startProcessing()(index.ts)。
  • startProcessing():迭代式调度主循环,控制并发并处理跳过队列(index.ts)。

1.2 Tarjan 算法模块

图论分析能力集中在 packages/service/core/workflow/utils/tarjan.ts,对外暴露四个函数:

  • findSCCs():Tarjan 算法找出所有强连通分量,返回nodeToSCC(节点→SCC ID)与sccSizes(SCC ID→分量大小)两个 Map(tarjan.ts)。
  • classifyEdgesByDFS():对全图执行一次 DFS,为每条边标注tree/back/forward/cross类型(tarjan.ts)。
  • isNodeInCycle():通过 SCC 分量大小判断节点是否位于循环中,sccSizes > 1即视为循环(tarjan.ts)。
  • getEdgeType():按source-target-sourceHandle三元组作为边的唯一键,从边类型表中查询边的类型(tarjan.ts)。

二、核心数据结构

2.1 边的状态机

运行时边的 schema 定义于 packages/global/core/workflow/type/edge.ts,在存储边(StoreEdgeItemTypesource/sourceHandle/target/targetHandle)基础上扩展出status字段:

type EdgeStatus = 'waiting' | 'active' | 'skipped';

三种状态的语义:

  • waiting:等待执行——源节点尚未完成,或该边对应的执行路径尚未走到;
  • active:已激活——源节点已执行完成,把结果"送达"了这条边;
  • skipped:已跳过——源节点被跳过(如分支未选中),该边所在路径被剪枝。

初始状态在 runtime/utils.ts 的storeEdges2RuntimeEdges中统一置为waiting,每次节点执行结束后由调度逻辑更新下游边状态(见 index.ts):命中skipHandleId的输出桩对应的边置为skipped,其余置为active

2.2 边的类型(图论视角)

type EdgeType = 'tree' | 'back' | 'forward' | 'cross';
  • tree:树边,DFS 树中首次发现目标节点的边;
  • back:回边,从后代指向当前 DFS 路径上祖先的边——即循环边;
  • forward:前向边,从祖先指向后代的非树边;
  • cross:跨边,连接不同 DFS 子树的边。

其中back边是循环判定的关键信号。

2.3 节点边分组

type NodeEdgeGroups = RuntimeEdgeItemType[][]; type NodeEdgeGroupsMap = Map<string, NodeEdgeGroups>;

每个节点的所有输入边会被划分成若干组:组内边是"且"关系(全部满足才能运行),组间是"或"关系(任意一组满足即可运行)。分组结构预构建一次、全流程复用,是状态判定的唯一依据。

三、核心算法

3.1 边分组算法流程

buildNodeEdgeGroupsMap的完整流程(源码注释见 index.ts):

1. 全局 DFS 边分类 └─> 识别回边(循环边) 2. Tarjan SCC 算法 └─> 找出所有强连通分量 └─> 判断节点是否在循环中 3. 为每个节点构建边分组 ├─> 分类边:回边 vs 非回边 ├─> 处理非回边 │ ├─> 节点在循环中 → 按 branchHandle 分组 │ └─> 节点不在循环中 → 所有非回边放在同一组 └─> 处理回边 └─> 按 branchHandle 分组

源码中isBranchNode只认三种节点类型:条件判断节点(ifElseNode)、意图分类节点(classifyQuestion)、用户选择节点(userSelect),见 index.ts。

3.2 分组策略

策略 1:节点不在循环中所有非回边放在同一组。这类边是"且"的关系,必须全部满足条件才能运行——典型代表是并行汇聚,多条并行支路全部完成后汇合节点才执行。

策略 2:节点在循环中非回边按branchHandle分组,回边也按branchHandle分组。不同组的边是"或"的关系,任意一组满足即可运行——这样既能保证首次从入口进入循环,也能保证循环体可以在回边激活后再次执行。

3.3 branchHandle 查找

findBranchHandle从边的源节点开始沿输入边向上回溯,找到第一个分支节点及其输出桩句柄(sourceHandle),以此作为该边归属的"分支路径"标识;若回溯不到任何分支节点,则返回默认值'common'。回溯过程中,若遇到分支节点则用其sourceHandle覆盖句柄,否则保持当前句柄继续向上(实现于 index.ts):

findBranchHandle(edge) { // 从边的源节点开始向上回溯 queue = [{ nodeId: edge.source, handle: edge.sourceHandle }] while (queue.length > 0) { { nodeId, handle } = queue.shift() // 如果当前节点是分支节点且有 handle,返回 handle if (isBranchNode(node) && handle) { return handle } // 继续向上回溯 for (inEdge of inEdges) { newHandle = isBranchNode(sourceNode) ? inEdge.sourceHandle : handle queue.push({ nodeId: inEdge.source, handle: newHandle }) } } return 'common' }

groupEdgesByBranch先为每条边计算branchHandle,再按句柄值聚合成组(index.ts)。

四、节点运行状态判断

节点只有三种运行状态:run(运行)、skip(跳过)、wait(等待)。判定逻辑getNodeRunStatus(index.ts)完全基于预构建的边分组:

getNodeRunStatus(node, nodeEdgeGroupsMap) { edgeGroups = nodeEdgeGroupsMap.get(node.nodeId) // 1. 没有输入边 → 入口节点,直接运行 if (!edgeGroups || edgeGroups.length === 0) { return 'run' } // 2. 检查是否可以运行(任意一组边满足条件) // 每组边内:至少有一个 active,且没有 waiting if (edgeGroups.some(group => group.some(edge => edge.status === 'active') && group.every(edge => edge.status !== 'waiting') )) { return 'run' } // 3. 检查是否跳过(所有组的边都是 skipped) if (edgeGroups.every(group => group.every(edge => edge.status === 'skipped') )) { return 'skip' } // 4. 否则等待 return 'wait' }

三条判定规则总结:

规则条件结果
运行条件任意一组边:至少一个active且没有waitingrun
跳过条件所有组的边全部为skippedskip
等待条件不满足以上两者wait

在调度主循环中(index.ts),run状态的节点会先把所有指向它的边重置为waiting再执行节点逻辑(nodeRunWithActive),skip状态的节点同样重置入边状态后执行跳过处理(nodeRunWithSkip),并消耗maxRunTimes预算(跳过节点扣减 0.1,正常运行扣减实际运行次数,见 index.ts 与 index.ts),maxRunTimes <= 0时整个工作流终止——这是防止死循环的兜底机制。

五、Tarjan 强连通分量算法

Tarjan 算法用于在有向图中找出所有强连通分量(Strongly Connected Components, SCC)。

强连通分量定义:在有向图中,若从节点 A 能到达节点 B,且从节点 B 也能到达节点 A,则 A、B 同属一个强连通分量。SCC 大小 > 1 即表示存在循环(自环节点自身构成大小为 1 的分量,但大小为 1 的分量不一定在循环中,因此循环判定只看size > 1)。

核心实现(tarjan.ts)使用lowLinkdiscoveryTime两个辅助表,以栈 + 迭代遍历的方式为每个节点分配 SCC ID:

function findSCCs(runtimeNodes, edgeIndex) { nodeToSCC = new Map() sccSizes = new Map() sccId = 0 stack = [] inStack = new Set() lowLink = new Map() discoveryTime = new Map() time = 0 function tarjan(nodeId) { // 初始化 discoveryTime.set(nodeId, time) lowLink.set(nodeId, time) time++ stack.push(nodeId) inStack.add(nodeId) // 遍历所有出边 for (edge of outEdges) { targetId = edge.target if (!discoveryTime.has(targetId)) { // 未访问过,递归访问 tarjan(targetId) lowLink.set(nodeId, min(lowLink.get(nodeId), lowLink.get(targetId))) } else if (inStack.has(targetId)) { // 在栈中,更新 lowLink lowLink.set(nodeId, min(lowLink.get(nodeId), discoveryTime.get(targetId))) } } // 如果是 SCC 的根节点 if (lowLink.get(nodeId) === discoveryTime.get(nodeId)) { sccNodes = [] do { w = stack.pop() inStack.delete(w) nodeToSCC.set(w, sccId) sccNodes.push(w) } while (w !== nodeId) sccSizes.set(sccId, sccNodes.length) sccId++ } } // 从所有未访问节点开始 for (node of runtimeNodes) { if (!discoveryTime.has(node.nodeId)) { tarjan(node.nodeId) } } return { nodeToSCC, sccSizes } }

六、DFS 边分类算法

classifyEdgesByDFS使用深度优先搜索对每条边进行分类(tarjan.ts),分类规则:

function classifyEdgesByDFS(runtimeNodes, edgeIndex) { edgeTypes = new Map() visited = new Set() inStack = new Set() discoveryTime = new Map() finishTime = new Map() time = 0 function dfs(nodeId) { visited.add(nodeId) inStack.add(nodeId) discoveryTime.set(nodeId, ++time) for (edge of outEdges) { targetId = edge.target if (!visited.has(targetId)) { // 未访问 → 树边 edgeTypes.set(edgeKey, 'tree') dfs(targetId) } else if (inStack.has(targetId)) { // 在当前路径上 → 回边(循环边) edgeTypes.set(edgeKey, 'back') } else if (discoveryTime.get(source) < discoveryTime.get(targetId)) { // 从祖先指向后代 → 前向边 edgeTypes.set(edgeKey, 'forward') } else { // 跨边 edgeTypes.set(edgeKey, 'cross') } } inStack.delete(nodeId) finishTime.set(nodeId, ++time) } // 从所有入口节点开始 DFS for (node of entryNodes) { if (!visited.has(node.nodeId)) { dfs(node.nodeId) } } return edgeTypes }

注意两点源码细节:

  • 边的唯一键是source-target-sourceHandle三元组(sourceHandle缺失时用'default'兜底),见 tarjan.ts,这意味着同一对节点之间允许存在多条不同输出桩的连线;
  • 入口节点定义为"没有输入边"的节点(tarjan.ts),DFS 结束后还会补遍历孤立节点(tarjan.ts),确保图中所有节点都被分类。

七、典型场景分析

设计文档给出了 5 个代表性拓扑,且每个场景在测试目录 packages/service/test/core/workflow/dispatch/checkNodeRunStatus 下都有对应的用例文件(base.test.tstoolcall.test.tssafe.test.tsboundary.test.tscase.test.ts)逐一验证。

7.1 简单分支汇聚

┌─ if ──→ B ──┐ start ──→ A ├──→ D └─ else ─→ C ──┘

边分组:D 节点只有一组[B→D, C→D]

运行逻辑

  • A 走 if 分支:B→D active, C→D skipped→ D 运行;
  • A 走 else 分支:B→D skipped, C→D active→ D 运行;
  • B 还在执行:B→D waiting, C→D skipped→ D 等待。

对应测试见 base.test.ts(场景 1:简单分支汇聚)。

7.2 简单循环

start ──→ A ──→ B ──→ C ──┐ ↑ | └────────────────┘

边分组:A 节点分为两组:组1[start→A]组2[C→A]C→A是回边)。

运行逻辑

  • 第一次执行:start→A active, C→A waiting→ A 运行(组 1 满足);
  • 循环执行:start→A skipped, C→A active→ A 运行(组 2 满足);
  • 两条边都 waiting:start→A waiting, C→A waiting→ A 等待。

对应测试见 base.test.ts(场景 2:简单循环)。

7.3 分支 + 循环

┌─ if ──→ B ──┐ start ──→ A ├──→ D ──┐ └─ else ─→ C ──┘ | ↑ | └──────────────────────┘

边分组:D 分为组1[B→D]组2[C→D](循环内按 branchHandle 分组);A 分为组1[start→A]组2[D→A]

运行逻辑

  • 第一次走 if 分支:B→D active, C→D skipped→ D 运行;
  • 第一次走 else 分支:B→D skipped, C→D active→ D 运行;
  • 循环回来:start→A skipped, D→A active→ A 运行。

对应测试见 base.test.ts(场景 3:分支 + 循环)。

7.4 并行汇聚(无分支节点)

start ──→ A ──→ C └──→ B ──→ C

边分组:C 只有一组[A→C, B→C](不在循环中,非回边合为同组,"且"关系)。

运行逻辑

  • A 和 B 都完成:A→C active, B→C active→ C 运行;
  • 只有 A 完成:A→C active, B→C waiting→ C 等待;
  • 只有 B 完成:A→C waiting, B→C active→ C 等待。

对应测试见 base.test.ts(场景 4:并行汇聚)。

7.5 工具调用场景

┌──selectedTools──→ Tool1 ──┐ start → Agent ─┤ ├──→ End └──────────────────────────→ ┘

边分组:Tool1 一组[Agent→Tool1 (selectedTools)];End 分为组1[Agent→End]组2[Tool1→End]

运行逻辑

  • Agent 调用 Tool1:Agent→Tool1 active→ Tool1 运行;
  • Agent 不调用工具:Agent→Tool1 skipped, Agent→End active→ End 运行;
  • Tool1 执行完成:Tool1→End active, Agent→End active→ End 运行(注意此时 Agent→End 也是 active,两组同时满足,体现分支与并行的混合语义)。

工具相关全部用例集中在 toolcall.test.ts,覆盖单工具、多工具并行、嵌套工具调用、工具与分支结合四类场景。

八、调度主循环:从队列到执行的完整链路

startProcessing(index.ts)是驱动整个工作流运转的迭代式主循环,源码注释(index.ts)概括了其设计:

  • 使用activeRunQueue记录待检查的节点(可能可以运行),并控制并发数量;
  • 每次添加新节点,以及节点运行结束后,均会执行一次processActiveNode检查;若未触发跳出条件,必定会继续从队列取节点处理;
  • 节点运行后,将 target 节点加入activeRunQueue等待下一轮处理。

循环内的关键控制点:

  1. 结束条件activeRunQueue与运行中的 Promise 集合都为空时,进入收尾——调试模式下若已无下一步运行节点且skipNodeQueue非空,则先处理跳过节点(index.ts)。
  2. 并发限制runningNodePromises.size >= this.maxConcurrency时用Promise.race等待最早完成的节点,从而在每次节点结束的瞬间立刻推进下一轮(index.ts)。
  3. 去重addActiveNode通过 Set 天然去重,同一节点不会在队列中出现两次(index.ts)。
  4. 边状态写入:节点执行完成后,依据skipHandleId批量更新所有下游边的active/skipped状态,并把下一批活跃/跳过节点收集起来(index.ts)。

九、性能优化

9.1 预构建边分组

  • 优化前:每次判断节点状态时都要重新计算边分组,时间复杂度 O(n × m)(n 为节点数,m 为边数);
  • 优化后:在WorkflowQueue初始化时一次性构建所有节点的边分组,运行时直接查Map,状态判定退化为 O(1)。

9.2 边索引

  • 优化前:每次查找节点的输入/输出边都要遍历所有边,时间复杂度 O(m);
  • 优化后:构建bySourcebyTarget两个 Map,查询降为 O(1)(EdgeIndex类型定义于 tarjan.ts)。

9.3 迭代替代递归

  • 优化前:使用递归处理节点队列,深层图可能导致调用栈溢出;
  • 优化后startProcessing使用while (true)迭代循环替代递归的processActiveNode,从根本上避免栈溢出问题。

十、测试覆盖

设计文档标注的测试文件路径为test/cases/global/core/workflow/dispatch/checkNodeRunStatus.test.ts,在当前仓库中对应实际路径为 packages/service/test/core/workflow/dispatch/checkNodeRunStatus(按文件拆分:basetoolcallsafeboundarycase)。

已覆盖场景(对应设计文档清单):

  1. 简单分支汇聚
  2. 简单循环
  3. 分支 + 循环
  4. 并行汇聚(无分支节点)
  5. 所有边都 skipped
  6. 多层分支嵌套
  7. 嵌套循环
  8. 多个独立循环汇聚
  9. 复杂有向有环图(多入口多循环)
  10. 自循环节点
  11. 用户工作流 - 多层循环回退
  12. 复杂分支与循环混合
  13. 多层嵌套循环退出
  14. 极度复杂多分支多循环交叉(部分场景)
  15. 工具调用 - 单工具场景
  16. 工具调用 - 多工具并行场景
  17. 工具调用 - 嵌套工具调用场景
  18. 工具调用 - 工具与分支结合场景

测试规模:设计文档记录总测试数 72、通过 72、失败 0(针对核心checkNodeRunStatus判定逻辑)。按当前仓库拆分后的文件统计,仅base.test.ts即含 72 个用例,加上toolcall(22)、safe(28)、boundary(11)、case(10),覆盖了从基础拓扑到边界与安全性(如死循环防护)的完整矩阵。

10.1 场景 14 的问题分析

问题:场景 14.7 测试失败——期望节点 F 在只有一条边 active 时等待,但实际返回run

原因:场景 14 包含 D→E 的交叉路径,导致 F 的两条输入边(D→F 和 E→F)被分成了不同的组。当 D→F active 时,第一组已满足条件,F 因此可以运行。

解决方案:删除场景 14.7 测试,理由有三:

  1. 场景 14 是极端复杂的测试场景,不应出现在实际工作流中;
  2. 在当前分组逻辑下,D→F 与 E→F 来自不同分支,它们是"或"的关系;
  3. 当 D→F active 时 F 可以运行,这本身符合分支逻辑的语义。

这一取舍体现了设计意图:运行时语义以"分支路径"为准,而非以"汇聚完整性"为准——来自不同分支的边天然是或关系,不做强制合流等待。

十一、设计原则

11.1 分支语义

  • "或"关系:来自不同分支的边是"或"的关系,任意一个分支满足条件即可运行——如 if-else 分支;
  • "且"关系:来自同一分支的边是"且"的关系,所有边都必须满足条件才能运行——如并行汇聚。

11.2 循环处理

  • 循环识别:使用 Tarjan SCC 算法识别循环,SCC 大小 > 1 即表示存在循环;
  • 循环边分组:回边(循环边)按branchHandle分组,不同循环路径的边分入不同组,保证循环体的重入能力。

11.3 应避免的复杂场景

  1. 跨分支的交叉路径(如 D→E);
  2. 多个循环出口(如 G→A 和 G→C);
  3. 过度嵌套的分支和循环。

原因:难以理解和维护、容易出现逻辑错误、性能开销大、用户体验差。这与场景 14.7 被删除的决策一脉相承——Runtime 在正确性与"图可理解性"之间做了明确取舍。

十二、未来优化方向

设计文档从三个方向展望了演进空间:

  1. 性能优化:并行执行优化(更智能的并发控制)、内存优化(减少中间状态存储)、缓存优化(缓存常用计算结果);
  2. 功能增强:更丰富的分支类型支持、更灵活的循环控制、更强大的错误处理;
  3. 可观测性:更详细的执行日志、更直观的执行可视化、更完善的性能监控。

十三、相关文件索引

核心代码

  • packages/service/core/workflow/dispatch/index.ts —WorkflowQueue类(边索引、边分组、状态判定、队列调度)
  • packages/service/core/workflow/utils/tarjan.ts — Tarjan SCC 算法、DFS 边分类
  • packages/global/core/workflow/runtime/type.ts — 运行时节点类型RuntimeNodeItemType与节点响应 schema
  • packages/global/core/workflow/runtime/utils.ts — 运行时工具函数(含storeEdges2RuntimeEdges初始化边状态)
  • packages/global/core/workflow/type/edge.ts — 存储边与运行时边的 schema 定义

测试文件

  • packages/service/test/core/workflow/dispatch/checkNodeRunStatus/base.test.ts — 基础拓扑(分支、循环、并行)状态判定测试
  • packages/service/test/core/workflow/dispatch/checkNodeRunStatus/toolcall.test.ts — 工具调用场景测试
  • packages/service/test/core/workflow/dispatch/checkNodeRunStatus/safe.test.ts — 边界与安全场景测试
  • packages/service/test/core/workflow/dispatch/checkNodeRunStatus/boundary.test.ts — 边界条件测试
  • packages/service/test/core/workflow/dispatch/checkNodeRunStatus/case.test.ts — 复杂组合场景测试

总结

FastGPT 工作流 Runtime 采用基于图论的设计:通过 Tarjan SCC 算法识别循环、DFS 边分类区分回边,再以"节点边分组"为桥梁,把任意复杂的有向图归约为"组间或、组内且"的判定模型,最终由WorkflowQueue的迭代式调度主循环统一驱动节点执行。分支、循环、并行、工具调用等场景在预构建边分组与三层状态判定(run / skip / wait)下均能获得确定且正确的执行语义。

通过预构建边分组、边索引、迭代替代递归等优化,Runtime 在保证正确性的同时具备良好的性能表现与栈安全特性,并配有覆盖 18 类场景的测试矩阵(当前仓库拆分后合计 143 个用例)作为正确性保障。未来可在并行执行、错误处理与可观测性方面继续演进,进一步释放复杂工作流编排的能力边界。

【免费下载链接】FastGPTFastGPT is a knowledge-based platform built on the LLMs, offers a comprehensive suite of out-of-the-box capabilities such as data processing, RAG retrieval, and visual AI workflow orchestration, letting you easily develop and deploy complex question-answering systems without the need for extensive setup or configuration.项目地址: https://gitcode.com/GitHub_Trending/fa/FastGPT

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询