跳到主要内容

调度内核:WorkflowQueue 怎么把一张图跑起来

30 秒导读: 用户在画布上连出的工作流是一张有向图(节点 + 边)。本章讲这张图怎么被跑起来——不是某个 LLM/检索节点内部怎么算,而是那台"发牌员":它决定谁先跑、谁后跑、哪条分支不走就把后面一串跳过、带环的循环怎么不死循环。核心就是一个类:WorkflowQueuepackages/service/core/workflow/dispatch/index.ts:343)。

本章定位:这是 FastGPT 工程含量最高的一章,是纯粹的执行引擎。图长什么样、节点里有哪些字段,看 01-workflow-data-model;一次对话怎么从 API 进来、SSE 怎么流式返回,看 02-chat-pipeline;单个 LLM 节点里的工具循环、检索算法,看 04-ai-nodes05-knowledge-base。这里只讲"把图跑起来"这一件事。


1. 这是什么(零基础也能懂)

一句话定义: WorkflowQueue 是一台图执行器——给它一堆节点和边,它按依赖顺序、带并发上限地把每个节点执行一遍,中途该跳过的分支跳过,该等的节点等齐上游再跑。

它要解决的问题,用大白话说:

你在 FastGPT 画布上画了这么一张流程:

开始 ──▶ 判断问题类型 ──(是"退货")──▶ 查订单 ──▶ 回复
└──(是"闲聊")──▶ 直接回复

现在用户问了句"我要退货"。引擎必须做到三件看似简单、其实很容易做错的事:

要做到说人话
顺序对"查订单"必须等"判断问题类型"跑完、且判定结果确实流向它,才能跑
分支干净既然走了"退货","闲聊 → 直接回复"这条整条都不能执行
不重复同一个节点,哪怕有好几条边指向它,也只跑一次

为什么不能直接"递归下去就完了"? 因为真实工作流有汇聚(一个节点等多个上游都到齐)、有分支(走了 A 就得干净地废掉 B)、有(loop 循环体会往回连边)、还有并发(几条独立支路想同时跑但又要限流)。天真的深度递归会栈爆、会重复执行、会在环里转不出来。

一句话直觉: 把它想成发牌桌——activeRunQueue 是"轮到谁出牌"的待办牌堆,发牌员一次只发几张(并发上限),每张牌打完会翻出下一批牌放回牌堆;打不成的支路进另一个"作废牌堆"(skipNodeQueue)单独清理。整局不靠一层套一层地喊话(递归),而是牌堆空了才散场

本节不碰代码细节。你只要记住:它是把一张图安全、有序、不重不漏地跑完的东西。


2. 顶层全景(它大概怎么转)

2.1 一张图看懂主循环

先给一句"怎么读这张图":中间那个 startProcessing 死循环是心脏,它反复从两个队列取节点交给 checkNodeCanRun 判定,判定结果又会往队列里塞新节点,直到两个队列都空 → 散场(resolve)。

入口节点 (isEntry)
│ addActiveNode()

┌───────────────────────────┐
│ activeRunQueue │ 待检查的节点(可能能跑)
│ (Set<id>) │
└───────────────────────────┘
│ startProcessing 循环取节点
│ (并发 ≤ maxConcurrency)

┌───────────────────────────┐
│ checkNodeCanRun │ 看入边状态 → run / skip / wait
└───────────────────────────┘
run │ skip │ wait │
▼ ▼ ▼
nodeRunWithActive nodeRunWithSkip (什么都不做,
跑真节点逻辑 只标记边为 等别的边到齐
callbackMap[type] skipped) 再被唤醒)
│ │
▼ ▼
nodeOutput:把出边标成 active / skipped,
收集下一批节点 → active 的塞回 activeRunQueue
skip 的塞进 skipNodeQueue


两个队列都空 & 无并发在跑 ──▶ resolve(this) 散场

2.2 部件一句话职责

部件干什么位置(dispatch/index.ts
WorkflowQueue整台引擎,一次 workflow 运行一个实例:343
activeRunQueue待"检查能否运行"的节点集合:370
skipNodeQueue待"执行跳过逻辑"的节点集合:371
maxConcurrency同时最多跑几个节点(默认 10):375
startProcessing非递归主循环,取节点、控并发、判散场:694
checkNodeCanRun单节点判定 + 执行 + 传播后续:1149
getNodeRunStatus只读判定:看入边给出 run/skip/wait:626
nodeOutput(闭包)跑完后更新出边状态、算出下一批节点:1222
callbackMap节点类型 → 具体 dispatch 处理器dispatch/constants.ts:38
Tarjan 工具识别环、判某节点是否在环里utils/tarjan.ts
WorkflowVariableState全局变量的运行态/存储态双形管理dispatch/utils/variables.ts:104

2.3 边状态:整个引擎的"信号灯"

这是全章最关键的一个概念,先建立它。 引擎不给节点直接打标记,而是给边打标记。每条边(RuntimeEdgeItemType)有一个 status,取值四种:

边状态含义
waiting上游还没跑到这条边(初始/等待态)
active上游跑完了,这条边"通电",数据可以往下游流
skipped上游决定这条分支不走,这条边"作废"
(无 waiting 也无 active 的中间态)getNodeRunStatus 综合判定

一个节点该不该跑,完全由它的入边状态决定——这就是下一节的主题。


3. 核心原理(逐个机制,由浅入深)

3.1 为什么用"队列 + 回调",不用递归

它要解决的小问题: 天真实现会写成"跑完一个节点,就递归调用下一个节点"。节点一多、链一长,调用栈会很深;分支/汇聚/环还会让同一节点被递归进入好几次。

源码作者自己写下了设计说明,值得原样对照(dispatch/index.ts:325-342WorkflowQueue 类头的块注释):

特点:
1. 可以控制一个 team 下,并发 run 的节点数量。
2. 每个节点,同时只会执行一个。一个节点不可能同时运行多次。
3. 都会返回 resolve,不存在 reject 状态。
方案:
- 采用回调的方式,避免深度递归。
- 使用 activeRunQueue 记录待运行检查的节点(可能可以运行),并控制并发数量。
- 每次添加新节点,以及节点运行结束后,均会执行一次 processActiveNode 方法……

三条不变式,记牢:

  1. 限流 —— 并发数受 maxConcurrency(默认 10,:375、构造函数 :391)约束。
  2. 同节点不并发 —— 一个节点同一时刻至多跑一份。
  3. 只 resolve,不 reject —— 单个节点内部出错也被吞成"跳过后续边 + 记录 error",不会把整台引擎 reject 掉。

这三条怎么落地? 分别对应三处代码:

(1)同节点去重靠 Set activeRunQueueSet<string>(存 nodeId),addActiveNode 进队前先查在不在(:681):

// 示意,非源码
addActiveNode(nodeId) {
if (this.activeRunQueue.has(nodeId)) return; // 已在队里就不重复加
this.activeRunQueue.add(nodeId);
if (!this.processingActive) this.startProcessing(); // 没在跑就启动主循环
}

真实实现见 dispatch/index.ts:681 addActiveNode。用 Set 天然去重,这就是"一个节点不会被排两次"的物理保证。

(2)"回调避免递归"落地为一个 while(true) 主循环。 注释里说的 processActiveNode 在当前版本已重构成迭代式 startProcessing:694,注释写明"迭代处理队列(替代递归的 processActiveNode)")。它的骨架:

// 示意,非源码;重点看"取节点—跑—回到循环"这个圈
async startProcessing() {
if (this.processingActive) return; // 防重复启动
this.processingActive = true;
const running = new Set(); // 正在跑的节点 Promise
while (true) {
// 散场条件:两个队列空 & 没有在跑的
if (activeRunQueue.size === 0 && running.size === 0) {
if (skipNodeQueue.size > 0 && !interactive) { await processSkipNodes(); continue; }
break;
}
// 并发满了 / 没有可取节点:等最快的一个跑完再回来
if (activeRunQueue.size === 0 || running.size >= maxConcurrency) {
await Promise.race(running); continue;
}
// 取一个节点,异步跑(不 await 它完成),继续循环去取下一个
const nodeId = activeRunQueue.keys().next().value;
activeRunQueue.delete(nodeId);
const p = checkNodeCanRun(node).finally(() => running.delete(p));
running.add(p);
}
this.resolve(this); // 散场
this.processingActive = false;
}

真实实现见 dispatch/index.ts:694 startProcessing关键手法: 取到节点后await 它跑完,而是把 Promise 塞进 running 集合继续循环去取下一个——这就是并发;只有当队列空或并发满时才 Promise.race(running) 等最先完成的那个(:731)。整局的结束由 this.resolve(this)finally 里触发(:748),对应 runWorkflownew Promise(resolve => …) 的那个 resolve(:1586)。

(3)"只 resolve 不 reject" 落地在节点执行处。 节点跑出错时,nodeRunWithActive 内部 catch 掉,把出边全标成 skipped 并记 error:909-931),而不是抛出去。所以下游要么被跳过、要么看到 error 字段,但引擎主循环不崩。

3.2 checkNodeCanRun:一个节点该跑、该跳、还是该等

它要解决的小问题: 从队列取出一个节点后,怎么知道它现在能不能跑?答案全在它的入边状态里。

判定逻辑在 getNodeRunStatus:626),三句话:

判定条件(看入边分组 edgeGroups
run(可以跑)没有入边(入口节点);存在某一组入边"至少一条 active 且没有 waiting"
skip(该跳过)所有组的所有入边都是 skipped
wait(再等等)其余情况——还有边悬而未决,先不动

对应源码(dispatch/index.ts:626 getNodeRunStatus,核心两段):

// check active(任意一组边满足条件即可运行)
// 每组边内: 至少有一个 active,且没有 waiting
if (edgeGroups.some(group =>
group.some(edge => edge.status === 'active') &&
group.every(edge => edge.status !== 'waiting'))) return 'run';

// check skip(所有组的边都是 skipped 才跳过)
if (edgeGroups.every(group => group.every(edge => edge.status === 'skipped'))) return 'skip';

return 'wait';

为什么入边要"分组"(edgeGroups)? 因为像 if-else 这种分支节点的多个出口,在下游汇聚时不能一刀切地"全 active 才跑"。分组由 buildNodeEdgeGroupsMap 在构造时一次性预计算:445),按 branchHandle 把边归组,再结合"节点是否在环里"决定分组策略。判定时直接查 nodeEdgeGroupsMap,不再实时遍历——这是性能优化("一次性计算,后续直接查询",:404)。

判定完做什么?checkNodeCanRun:1149)里:

  • run → 先把该节点所有入边置 waiting:1301-1305,防止本轮被重复触发),做余额检查,再调 nodeRunWithActive 真跑(:1318)。
  • skip 且没跳过过 → 入边置 waitingmaxRunTimes -= 0.1(跳过也计一点点消耗防死循环,:1327),调 nodeRunWithSkip:1331)。
  • waitnodeRunResult 为空,直接 return:1338),节点静静躺着,等未来某条边变 active 时它会被重新 addActiveNode 唤醒。

3.3 节点跑完:更新边状态 + 收集下一批(nodeOutput)

跑完一个节点,引擎怎么知道下一步是谁? 靠闭包 nodeOutput:1222)。它做三件事:

  1. 把节点输出写回 node.outputs[].value:1240),供下游读。
  2. 更新出边状态:节点返回一个 skipHandleId 列表(哪些出口该被跳过)。命中的出边标 skipped,其余标 active:1253-1260):
targetEdges.forEach((edge) => {
if (skipHandleId.includes(edge.sourceHandle)) {
edge.status = 'skipped';
} else {
edge.status = 'active';
}
});
  1. 收集下一批节点并去重:遍历出边,active 的目标进 nextStepActiveNodesskipped 的进 nextStepSkipNodes,用 Map 去重(:1263-1275)。

回到 checkNodeCanRun 尾部:nextStepSkipNodes 逐个 addSkipNode:1382),nextStepActiveNodes 逐个 addActiveNode:1419)。这就是"翻牌" ——跑完一张牌,翻出下一批放回牌堆,主循环下一轮自然取到。

3.4 跳过传播(skip propagation):让废掉的分支干净地"死透"

它要解决的小问题: if-else 走了"true"分支,"false"那一整条支路必须全部不执行。但下游节点可能被别的活支路汇聚,不能见 skip 就无脑往下传。

机制:skip 也是"节点",也要被处理。 被判 skip 的节点进 skipNodeQueueMap:764 addSkipNode),nodeRunWithSkip:1115)不跑真逻辑,只返回"我的所有出边都进 skipHandleId":

// 示意,非源码
nodeRunWithSkip(node) {
const targetEdges = edges.filter(e => e.source === node.nodeId);
return { runStatus: 'skip',
skipHandleId: targetEdges.map(e => e.sourceHandle) }; // 出边全标 skipped
}

于是它的出边全被 nodeOutput 标成 skipped,下游节点被 addSkipNode 排进 skip 队列——skip 就这样沿边一路传下去。而下游若还有别的 active 入边,getNodeRunStatus 会因为"不是所有边都 skipped"而判 runwait自动止住误传

为什么 skip 队列和 active 队列分开、且 skip 每次只处理一个? 看主循环:只有当 activeRunQueue 空了、没有并发在跑、也没有待响应的交互时,才去 processSkipNodes:719:775),而且每次只清一个就 continue 回主循环:775 注释"每次只处理一个,然后返回主循环检查 active")。这保证能跑的节点优先跑,跳过是"收尾清理",避免 skip 抢跑把本该 active 的节点误判成 skip。

环场景下的防递归: 分支节点的 skip 边可能往回连(环),无限递归 skip。代码用 skippedNodeIdList 记录已跳过的节点(checkNodeCanRun 的第二参数,:1149),并对分支节点特判(:1365-1372 块注释 + :1371 skippedNodeIdList.add):

特殊情况:
通过 skipEdges 可以判断是运行了分支节点。
由于分支节点,可能会实现递归调用(skip 连线往前递归)
需要把分支节点也加入到已跳过的记录里,可以保证递归 skip 运行时,
至多只会传递到当前分支节点,不会影响分支后的内容。

skippedNodeIdList 随传播路径累积,某节点已在集合里就不再重复跳(:1320!skippedNodeIdList.has(node.nodeId) 守卫)——这是 skip 不会在环里打转的关键闸门

3.5 环的处理:Tarjan SCC 让 loop 也能安全调度

它要解决的小问题: loop 循环体天然带回边(循环末尾连回开头)。如果引擎把回边和普通前向边一视同仁,getNodeRunStatus 会永远等那条"还没跑的回边",死锁。

思路:先把图分析清楚,认出哪些边是"回边"、哪些节点"在环里",判定时区别对待。 这在构造 WorkflowQueue一次性算好buildNodeEdgeGroupsMap:445),用 utils/tarjan.ts 三个函数:

函数干什么位置
classifyEdgesByDFS一遍 DFS 给每条边分类:tree / back(回边=环边) / forward / crosstarjan.ts:100
findSCCsTarjan 算法找出所有强连通分量(SCC),大小 > 1 即成环tarjan.ts:20
isNodeInCycle查某节点所在 SCC 大小是否 > 1,判它是否在环里tarjan.ts:79
getEdgeTypesource-target-handle 键取回该边被分类成什么tarjan.ts:170

术语一句话点破:

  • 强连通分量(SCC)——一组节点,两两之间都能互相到达。SCC 大小 > 1 就意味着这组节点构成了
  • 回边(back edge)——DFS 时指向"当前路径上还没出栈的祖先"的边,就是把图变成环的那条边(tarjan.ts:128-130)。

分组时怎么用(buildNodeEdgeGroupsMap:445 起):

对每个目标节点:
targetInCycle = isNodeInCycle(node) # 它在环里吗?
把入边分成 backEdges / nonBackEdges # getEdgeType 判回边
├─ 非回边:若节点在环里 → 按 branchHandle 分组;否则整组
└─ 回边:单独处理(不让"还没跑的回边"阻塞首次进入)

效果: 环里的节点第一次能靠"非回边 active"正常进入(不被尚未通电的回边卡住),循环体每跑一圈通过回边把状态送回入口,loop 就能一圈圈转,又不会因回边永远 waiting 而死锁。注意: 防无限循环的兜底是 maxRunTimes——每跑一个节点 maxRunTimes -= runTimes:1167),到 0 就强制停(:1290),loop 的实际迭代上限还叠加 loopRun 自己的次数控制(见 3.7)。

3.6 callbackMap:节点类型 → 处理器的分派表

主循环只负责调度,"节点里具体干啥"外包给分派表。 nodeRunWithActive 里,真正执行是这一行(:841):

const result = await callbackMap[node.flowNodeType](dispatchData);

callbackMap 是一张 节点类型 → dispatch 函数 的静态表(dispatch/constants.ts:38)。节选:

节点类型处理器属于哪章
workflowStartdispatchWorkflowStart本章(入口)
chatNode / agentdispatchChatCompletion / dispatchRunAgent04-ai-nodes
datasetSearchNodedispatchDatasetSearch05-knowledge-base
ifElseNode / classifyQuestion / userSelect分支类,产出 skipHandleId本章(分支/交互)
loopRun / parallelRundispatchLoopRun / dispatchParallelRun本章 3.7(嵌套子运行)
formInput / userSelectdispatchFormInput / dispatchUserSelect本章 3.8(交互)
emptyNode / comment / toolSet() => Promise.resolve() 空实现——

引擎与节点解耦的意义: 加一种新节点,只要写个 dispatch 函数、在这张表登记一行,调度内核一行不用改。dispatchData 把引擎的运行态(变量、histories、runtimeNodes/Edges、SSE writer 等)打包传给处理器(:822-838)。

3.7 loop / parallel:节点里再开一台引擎(嵌套子运行)

loop 和 parallel 不是引擎内建原语,而是"会递归调用引擎自己"的普通节点。 dispatchLoopRun 对数组每一项,克隆一份子图节点/边、把循环体入口标 isEntry = true、再调 runWorkflow 开一台子引擎loopRun/runLoopRun.ts:171 置 isEntry、:185runWorkflow):

// 示意,非源码:loop 每一轮开一次子运行
for (const [index, item] of list.entries()) {
injectLoopRunStart({ nodes, item, index }); // 把当前项喂给循环体入口
const response = await runWorkflow({
...props,
runtimeNodes: isolatedNodes, // 只含循环体的子图
runtimeEdges: isolatedEdges,
}); // ← 又是一整台 WorkflowQueue
}

parallel 同理,对每条并行支路 runWorkflowparallelRun/runParallelRun.ts:96)。深度有护栏: runWorkflow 每进一层 workflowDispatchDeep + 1,超过 20 直接返回空结果(:1499:1506),防止嵌套无底洞。

一句话: 主图跑到 loopRun 节点,就在这个节点内部递归开一台完整的 WorkflowQueue 把子图跑 N 遍,跑完把结果收集回来当作该节点的输出,主图继续。所以整套调度逻辑对循环体自动复用。

3.8 交互暂停与恢复:把一台跑到一半的引擎"冻"下来

它要解决的小问题: 遇到"表单输入""让用户选一个"这类节点,工作流得停下来等人,把当前状态存住,用户回答后再从断点续跑。

暂停: 交互节点(userSelect / formInput)跑完会返回一个 interactive 响应。checkNodeCanRun 检测到它就不再把下游 active 节点入队,而是把交互信息记到 nodeInteractiveResponse:1396-1414)后直接 return。主循环里也特判:有待响应的交互时不清 skip 队列:719!this.nodeInteractiveResponse),让 skip 状态原样保留。于是引擎自然停下、resolve 散场。

存档: handleInteractiveResult:1426)把现场打包成可持久化的 WorkflowInteractiveResponseType:收集所有节点当前 outputs、把 skipNodeQueue 序列化、把入口前的边全标成 active 保证下次一定能续上(:1461 注释"入口前面的边全部激活,保证下次进来一定能执行")、记录 entryNodeIds。这份存档随 SSE 发给前端保存。

恢复: 用户回答后再次进入 runWorkflow,靠 isEntry 标记定位断点节点(:1568 entryNodes = data.runtimeNodes.filter(item => item.isEntry))。注意"重置 entry"时特意跳过 userSelect / formInput / toolCall:1571-1578),因为这几类要靠 isEntry 当续跑入口。defaultSkipNodeQueue 从上次存档恢复(:1589),lastInteractive 带着上次答案喂回对应节点(:801-828)。断点续跑就这样接上。

3.9 debug 模式:一次只走一步

debug 模式让工作流单步执行。 普通模式下 nextStepActiveNodes 会立刻 addActiveNode 继续跑;debug 模式改成只记不跑——存进 debugNextStepRunNodes:1422),主循环发现"没有下一个激活节点了"就进入"即将结束"态、开始处理 skip(:707-714),最后 getDebugResponse:1480)把"下一步该跑哪些节点(debugNextStepRunNodes)+ 当前边状态 + skip 队列"整包返回前端。用户点"下一步"时,这批节点作为新入口再跑一轮。单步调试 = 每次只放行 debugNextStepRunNodes 这一层。


4. 深入实现(关键数据结构与时序)

4.1 变量状态:WorkflowVariableState

节点之间靠全局变量传值,统一由 WorkflowVariableStatedispatch/utils/variables.ts:104)管。它的设计要点(文件头注释 :20-33):同一份变量维护两种形态——

形态用途
storeValue可持久化落库、可回传前端继续编辑的值
runtimeValue节点运行时直接消费的值(如 file 变量转成 URL 数组)

统一入口:读走 get():153,返回 runtimeValue),写走 set():170),最终存档走 toStoreRecord():231,只导出非 runtimeOnly 的 storeValue)。特殊变量(file/password/runtimeOnly)的转换全集中在这里,避免散落各节点。子运行拿到 parent 传下的 file URL,还能通过 sourceVariableState 找回原始 store metadata(:165)。

4.2 一次节点执行的时序(nodeRunWithActive)

nodeRunWithActive:787)是"真跑一个节点"的全过程,主要步骤:

1. 算 nodeResponseId(续跑时复用上次的 id,:798-803)
2. showStatus 节点 → SSE 推 "running" 状态(:806)
3. getWorkflowNodeRunParams:从变量态 + 上游 outputs 拼本节点入参(:817)
4. callbackMap[type](dispatchData) ← 真正执行,:841
├─ 抛错 & 非 catchError → 出边全 skipped + 记 error(:846-861)
├─ 返回 error 字段 → 同上统一成 nodeResponse.error(:843-870)
└─ catchError 节点 → 只跳非 error 出口,保留 error 分支(:872-885)
5. 落库 nodeResponse(nodeResponseWriter.record,:970)
6. apiVersion v2 → SSE 推 flowNodeResponse(:986)
7. 返回 { runStatus:'run', nodeResponseId, result }

统一错误约定值得留意: 无论处理器是 throw 还是返回 { error },最终都被归一成 nodeResponse.error:843 起注释解释:为让 runLoopRun/parallelRun 的失败检测和 OTel span 状态在两条失败路径上看到一致的 .error)。

4.3 入队一次、检查一次的完整闭环

把 3.1~3.3 串起来,一个节点的生命周期:

addActiveNode(id) // 进 activeRunQueue(Set 去重)
↓ startProcessing 取出
checkNodeCanRun(node)
↓ getNodeRunStatus 看入边
├─ run → 入边置 waiting → nodeRunWithActive → callbackMap 执行
├─ skip → nodeRunWithSkip(只标出边)
└─ wait → return(躺平,等未来某条边 active 再被 addActiveNode 唤醒)
↓(run/skip 都会)
nodeOutput:出边标 active/skipped,算 next 节点
├─ nextActive → addActiveNode(回到顶部,闭环)
└─ nextSkip → addSkipNode(进 skipNodeQueue)

5. 巧妙之处(可借鉴的技术)

  1. 状态挂在边上,不挂在节点上。 节点该跑/该跳/该等,全由入边 status 组合推出(getNodeRunStatus:626)。这让"分支""汇聚""跳过传播"用同一套规则表达,无需给节点加复杂状态机。

  2. "能跑优先,跳过收尾"的双队列。 active 队列跑干净了才清 skip 队列(:719),且 skip 每次只清一个就回主循环(:775)。避免跳过逻辑抢跑,把本该 active 的节点误判成 skip。

  3. 图分析一次性预计算,判定时只查表。 Tarjan SCC + 边分类在构造时算好塞进 nodeEdgeGroupsMap:445),运行时 getNodeRunStatus 直接查(:634),把"每次判定都遍历图"降成 O(1) 查表。

  4. skippedNodeIdList 累积集当环闸门。 分支节点跳过时把自己加进已跳过集(:1371),保证 skip 沿环回溯"至多传到当前分支节点"(:1365 块注释),不误伤分支之后的内容。

  5. loop/parallel = 递归开引擎,而非内建原语。 循环/并行只是会调 runWorkflow 的普通节点(runLoopRun.ts:185),配 workflowDispatchDeep > 20 深度护栏(:1499)。调度逻辑对嵌套天然复用。

  6. 只 resolve 不 reject。 单节点错误被吞成"跳边 + 记 error"(:909),整台引擎永不因单点失败崩掉,错误以数据形式流向下游或前端。


6. 边界与局限

  • 同一时刻只允许一个交互节点:1408 注释 "only one interactive node is allowed at the same time";仅 paymentPause 例外,可累积多个 entryNodeId,:1400)。多个交互并发暂停不被支持。
  • 嵌套深度硬上限 20:1499),超过直接返回空结果,深层嵌套工作流会被截断。
  • 防死循环靠 maxRunTimes 配额:1167:1290),不是靠"检测到无进展"。配额是运行次数预算,跳过也扣 0.1(:1327)。
  • 调度不理解节点语义——它只看边状态和 skipHandleId,节点内部逻辑(LLM 调用、检索)完全是黑盒,由 callbackMap 处理器负责。本章不涉及这些内部算法。
  • 停止/中断分 v1(客户端 abort)与 v2(Redis 停止标记轮询,:238)两套,checkIsStopping 每步都查(:1283)。

7. 横向对比

同 shelf 的 agent 框架里,"图执行"是常见能力,取舍不同:

维度FastGPT WorkflowQueue典型代码优先框架(如 LangGraph 类)
图来源可视化画布产出的节点/边代码里声明的图
调度单位边状态驱动的活跃队列 + 跳过队列多为 super-step / 消息传递
环支持Tarjan SCC 预分析 + 回边特判靠显式条件边 + 步数上限
失败模型只 resolve,错误变边状态多为异常上抛/中断

FastGPT 的取向是面向可视化编排:把"哪条分支不走""汇聚等谁"这类画布语义,用边状态一套规则吃掉,让非程序员也能连出可靠运行的图。


8. 代码地图(导航索引)

主题文件路径符号名
引擎主体 + 设计说明注释packages/service/core/workflow/dispatch/index.ts:343WorkflowQueue
非递归主循环 / 并发控制 / 散场dispatch/index.ts:694startProcessing
单节点判定 + 执行 + 传播dispatch/index.ts:1149checkNodeCanRun
run/skip/wait 判定(看入边)dispatch/index.ts:626getNodeRunStatus
出边状态更新 + 收集下一批dispatch/index.ts:1222nodeOutput(闭包)
真正执行节点 / 调 callbackMapdispatch/index.ts:787 / :841nodeRunWithActive
跳过节点(只标出边)dispatch/index.ts:1115nodeRunWithSkip
入活跃队列(Set 去重)dispatch/index.ts:681addActiveNode
入跳过队列(Map 累积)dispatch/index.ts:764 / :775addSkipNode / processSkipNodes
交互暂停存档dispatch/index.ts:1426handleInteractiveResult
debug 单步返回dispatch/index.ts:1480getDebugResponse
入口装配 + 深度护栏dispatch/index.ts:1499runWorkflow
环分析:SCCpackages/service/core/workflow/utils/tarjan.ts:20findSCCs
环分析:边分类utils/tarjan.ts:100classifyEdgesByDFS
环分析:是否在环utils/tarjan.ts:79 / :170isNodeInCycle / getEdgeType
节点类型 → 处理器分派表dispatch/constants.ts:38callbackMap
全局变量双形态管理dispatch/utils/variables.ts:104WorkflowVariableState
loop 嵌套子运行dispatch/loopRun/runLoopRun.ts:185dispatchLoopRunrunWorkflow
parallel 嵌套子运行dispatch/parallelRun/runParallelRun.ts:96dispatchParallelRunrunWorkflow