AgentFlow V2 引擎:队列驱动的有状态解释器
30 秒导读: 这是 Flowise 的第二套执行引擎,只跑
type=AGENTFLOW的图。它不像 经典引擎那样"把整张图拓扑排序、一次性构建成一条 LangChain 流水线再 run 一次", 而是像解释器一样,从起点节点开始、一个节点一个节点地执行:执行完一个节点,看它的输出决定下一步 把哪些节点推进队列。因为控制权始终握在引擎手里,所以循环、条件分支、人在环暂停/恢复、子流嵌套 这些"经典 DAG 表达不了"的东西,在这里都成了自然而然的事。
1. 为什么要新引擎(先讲清动机)
1.1 经典引擎的世界观:构建一次,终点跑一次
经典引擎的核心动作是拓扑构建:它把整张图按依赖排好序,从叶子往根依次
initNode,把每个节点实例化并用边把它们接成一条 LangChain 链,最后在终点节点上 run() 一次,
数据顺着这条预先接好的管道流到底。
这套世界观的隐含前提是:图是有向无环的(DAG),数据只朝一个方向流一次。 它非常适合 "检索 → 提示词拼装 → LLM → 输出"这类直线型 RAG 管道。
1.2 DAG 表达不了的四件事
一旦 agent 编排的需求上来,DAG 就捉襟见肘。下面四类需求,经典引擎都很难自然表达:
| 需求 | 为什么 DAG 难 |
|---|---|
| 循环(反复改写直到满意) | 边不能指回上游,否则不再是"无环",拓扑排序直接失败 |
| 条件分支(走 A 还是走 B) | 管道一旦接好就固定了,无法"运行到一半才决定不走某条边" |
| 人在环(停下等人点"通过/拒绝") | 一次性 run 到底的模型里,没有"暂停、把状态存下来、以后再从这里续"的位置 |
| 长时运行 / 子流嵌套 | 单次 run 的调用栈里,塞不下"跑另一个完整 agentflow 再回来" |
1.3 V2 的答案:把"构建"换成"解释执行"
AgentFlow V2 换了世界观:引擎自己拿着一个待执行队列,一个节点一个节点地解释执行。
- 执行完一个节点,引擎当场看它的输出,动态决定下一步把哪些子节点推进队列(走哪条边)。
- 想循环?把上游节点再推一次队列即可(用计数器防止无限循环)。
- 想暂停等人?执行到人在环节点就把整个执行状态存进数据库并 return,人点了按钮再从存档续跑。
一句话对照:
经典引擎: 图 ──拓扑构建──▶ 一条 LangChain 链 ──run 一次──▶ 结果
(管道预先焊死,数据流一遍)
V2 引擎: 图 ──▶ [起点入队] ──▶ 循环{ 出队一个节点 → 执行 → 看输出决定谁入队 } ──▶ 结果
(管道边跑边铺,控制权始终在引擎手里)
引擎主体全部在一个文件里:packages/server/src/utils/buildAgentflow.ts。下面逐层拆开。
2. 顶层全景(引擎大概怎么转)
整个引擎就是一个 while 循环 + 三个状态容器。先看这张图,再逐块讲。
executeAgentFlow() ← 入口(buildAgentflow.ts:1541)
│
┌───────────────────┴────────────────────┐
│ 初始化三个状态容器 (buildAgentflow.ts:1650-1652)│
│ nodeExecutionQueue : 待执行节点队列 │
│ waitingNodes(Map) : 还没凑齐输入的节点 │
│ loopCounts(Map) : 每个循环节点跑了几圈 │
└───────────────────┬────────────────────┘
│ 把起点(startAgentflow)入队
▼
┌──────────────────────────────────────────┐
│ while(队列非空 && status=INPROGRESS) (2020)│
│ │
│ ① shift 出队一个节点 │
│ ② executeNode() 真正跑它 (1051) │
│ ③ 把这步结果 push 进 agentFlowExecutedData │
│ ④ processNodeOutputs() 决定下一步 (847) │
│ ├ 条件未命中的分支 → 剪掉 (799) │
│ ├ 子节点凑齐输入 → 入队 │
│ └ loop 节点 → 把目标节点再入队 │
│ ⑤ 人在环 / 等待工具确认 → shouldStop, break │
└──────────────────────────────────────────┘
│
▼
updateExecution() 落库 + 组装响应 (2230)
怎么读这张图: 从上到下是一次调用的生命周期;中间那个方框是引擎的心脏——一个反复"出队→执行→ 决定下一步"的循环。左侧三个容器是它全程依赖的"记事本"。
各部件一句话职责:
| 部件 | 干什么 | 位置(buildAgentflow.ts) |
|---|---|---|
executeAgentFlow | 引擎入口:初始化、跑主循环、落库、组装响应 | :1541 |
nodeExecutionQueue | 待执行节点的 FIFO 队列,主循环从这里取活 | :1650 |
waitingNodes | 多入边节点的"等待室":还没凑齐上游输入就先待在这 | :1651 |
loopCounts | 每个 loop 节点已经绕了几圈,用来防死循环 | :1652 |
executeNode | 执行单个节点:解析变量 → 调节点 run() → 处理人在环/迭代 | :1051 |
processNodeOutputs | 看节点输出,把该走的子节点入队、该剪的分支剪掉 | :847 |
IAgentFlowRuntime | 跨节点流转的运行态(state / chatHistory / form / webhook) | :97 |
addExecution / updateExecution | 把每一步执行数据写进 Execution 表,支持回看与恢复 | :184 / :212 |
3. 三个状态容器(引擎的"记事本")
主循环之所以能表达复杂控制流,全靠这三个数据结构。它们在 executeAgentFlow 开头一起初始化
(buildAgentflow.ts:1650-1652)。
3.1 nodeExecutionQueue —— 待办队列
一个数组,元素是 INodeQueue(buildAgentflow.ts:75):
| 字段 | 含义 |
|---|---|
nodeId | 该执行哪个节点 |
data | 合并后的上游输入(喂给节点) |
inputs | 各上游来源分开的原始输入({sourceNodeId: result}) |
主循环每轮 shift() 取队首,跑完再把新就绪的子节点 push 进来。队列空了,执行就结束。
3.2 waitingNodes —— 等待室
不是每个节点被"点到"就能立刻跑。有多条入边的节点,得等所有该到的上游都到齐才行。
waitingNodes 是个 Map<nodeId, IWaitingNode>,IWaitingNode(buildAgentflow.ts:67)记着:
expectedInputs:这个节点必须等到的上游集合。receivedInputs:已经到了的上游 → 结果。isConditional/conditionalGroups:如果某些入边来自条件分支,它们归成一组——同组 只要到一个就算数(因为条件分支本来就只会走其中一条)。
判断"能不能跑了"的逻辑在 hasReceivedRequiredInputs(buildAgentflow.ts:770):必需输入
全到齐、且每个条件组至少来一个,才返回 true。
3.3 loopCounts —— 圈数计数器
Map<nodeId, number>,记每个 loopAgentflow 节点已经循环了多少圈。每次循环回跳前 +1,
到达上限就停——防止无限循环的兜底就靠它(详 见 §5.3)。上限常量:
// buildAgentflow.ts:174 —— 允许用环境变量覆盖,默认 10 圈
const MAX_LOOP_COUNT = process.env.MAX_LOOP_COUNT ? parseInt(process.env.MAX_LOOP_COUNT) : 10
4. 执行主循环(executeAgentFlow)
这是引擎的骨架。入口 buildAgentflow.ts:1541,主循环 buildAgentflow.ts:2020。
4.1 起点从哪来
图先被 constructGraphs 拆成正向 graph 和 nodeDependencies(每个节点的入边数),
再各求一份反向图(buildAgentflow.ts:1595-1596)。源码注释直接给了这两个结构的形态
(buildAgentflow.ts:1624-1642):
graph { nodeDependencies {
startAgentflow_0: [ 'conditionAgentflow_0' ], startAgentflow_0: 0, ← 入度0 = 起点
conditionAgentflow_0: [ 'llmAgentflow_0', conditionAgentflow_0: 1,
'llmAgentflow_1' ], llmAgentflow_0: 1,
llmAgentflow_0: [ 'llmAgentflow_2' ], llmAgentflow_1: 1,
llmAgentflow_1: [ 'llmAgentflow_2' ], llmAgentflow_2: 2 ← 两条入边
llmAgentflow_2: []
} }
graph[nodeId]告诉引擎"跑完这个节点,候选下一步是谁"。nodeDependencies里入度为 0 的就是起点;getStartingNode(server/src/utils/index.ts:213) 从中挑出起点,checkForMultipleStartNodes保证非递归调用只有一个起点 (buildAgentflow.ts:1509)。起点入队在buildAgentflow.ts:1904。llmAgentflow_2的nodeDependencies=2正是它要在waitingNodes里等两个上游的原因。
4.2 循环体一轮做什么
主循环体(buildAgentflow.ts:2020-2211)每一轮:
- 出队一个节点,找到它的
reactFlowNode;跳过便签节点(:2032-2036)。 - 执行
executeNode(...)(:2048)——真正调用节点实现,见 §5。 - 若返回
shouldStop(人在环 / 等待工具确认),置status=STOPPED并break(:2102)。 - 否则把这步结果 push 进
agentFlowExecutedData,标FINISHED(:2110)。 - 把运行态回写:节点返回的
state/chatHistory/form就地更新到agentflowRuntime(:2127-2137,见 §6)。 - 分流
processNodeOutputs(...)(:2144)——决定谁入队、谁被剪、要不要循环回跳,见 §5。
有两道防线:MAX_ITERATIONS(默认 1000,:1912/2028)防止主循环本身跑飞;每轮开头查
abortController.signal.aborted(:2041)以支持用户中止。
5. 单节点执行与分流
主循环的两个关键子程序:executeNode(跑一个节点)和 processNodeOutputs(决定下一步)。
5.1 executeNode —— 跑一个节点
buildAgentflow.ts:1051。核心步骤(都在 try 块内,出错就抛给主循环的 catch 统一处理):
- 动态加载节点实现:按
componentNodes[name].filePathimport节点类并new(:1113-1115)——所以每个节点跑的是packages/components/nodes/agentflow/*里各自的run()。 - 解析变量:调
resolveVariables(...)(:1169)把{{ ... }}引用替换成真值(见 §6.2)。 - 组装
runParams并newNodeInstance.run(reactFlowNodeData, finalInput, runParams)(:1259)。 - 迭代节点特判:若是
iterationAgentflow且输出带数组,就对每个元素递归调用executeAgentFlow跑子块(:1262-1406,见 §7.4)。 - 人在环特判:若是
humanInputAgentflow且本轮没带humanInput,就构造带 "通过/拒绝"按钮的humanInputAction,把这步标STOPPED落进执行数据,返回shouldStop:true(:1409-1452)。agentAgentflow在"用工具前要人批准"时同理(:1456-1500)。
5.2 processNodeOutputs —— 决定下一步
buildAgentflow.ts:847。这是"边跑边铺管道"的地方:
processNodeOutputs(当前节点, 结果):
childNodeIds = graph[nodeId] # 候选下一步
ignoreNodeIds = determineNodesToIgnore(...) # 条件没命中要剪的分支 (799)
for childId in childNodeIds:
if childId 在 ignore 里: continue # 剪枝
waitingNode = waitingNodes[childId] 或新建 # 进等待室 (885-891)
waitingNode.receivedInputs[nodeId] = 结果 # 记下"我这条边到了"
if hasReceivedRequiredInputs(waitingNode): # 该到的都到齐了?(897)
从等待室删除, 合并输入, 入队 # 就绪 → push 队列 (899-904)
else:
继续待在等待室 # 还差上游, 先等着
多入边节点的输入用 combineNodeInputs(buildAgentflow.ts:971)按来源合并成一个对象再喂下去。
5.3 determineNodesToIgnore —— 条件分支剪枝
buildAgentflow.ts:799。条件节点如何"不走某条边"就靠它。
- 只对决策节点(
conditionAgentflow/conditionAgentAgentflow/humanInputAgentflow)生效。 - 节点输出里带一组
conditions,每个有isFulfilled布尔。没命中的条件对应的那条出边 (sourceHandle = ${nodeId}-output-${index})的目标节点,就被加进ignoreNodeIds剪掉 (:824-838)。 - 兜底:如果一个条件都没命中,就把最后一条当 else 分支强制置为命中,保证至少走一条
(
:818-822)——避免"全剪光、流程凭空断掉"。
5.4 循环回跳
还是在 processNodeOutputs 里(buildAgentflow.ts:918):如果刚跑完的是 loopAgentflow 且输出带
nodeID(要跳回的目标),就:
loopCount = loopCounts[nodeId] + 1
if loopCount < maxLoop: # 没到上限
loopCounts[nodeId] = loopCount
把 result.output.nodeID 再 push 进队列 # ← 关键:目标节点重新入队 = 循环
清掉 humanInput 防止被重复消费
else: # 到上限
流式发一条 fallbackMessage, 停止循环
这就是循环的全部秘密:不修改图、不加环,只是"把一个已经跑过的节点再推一次队列"。计数器守住上限。
6. 运行态:数据怎么在节点间流转
经典引擎里数据顺着焊死的管道流;V2 里节点之间不直接连数据——它们共享一本"运行态记事本"。
6.1 IAgentFlowRuntime
buildAgentflow.ts:97。四个字段,全程被主循环携带、被节点读写:
| 字段 | 装什么 | 谁写 |
|---|---|---|
state | 跨节点共享的键值状态(工作流的"全局变量") | 节点返回 state 时回写(:2127) |
chatHistory | 本次执行内累积的对话轮次 | 节点返回 chatHistory 时追加(:2131) |
form | 表单输入(formInput 起点持久化在 session 内) | :2135、:1734 |
webhook | 触发本次执行的 webhook 原始载荷 | :1770 |
它在 executeAgentFlow 开头初始化成空壳(:1655-1660),可能从上一次执行恢复(持久化 state、
表单值、webhook 重放,见 :1683-1757)。
6.2 resolveVariables 如何读运行态
buildAgentflow.ts:241。节点输入里写的 {{ ... }} 模板,跑之前由它替换成真值。它认得一整套前缀,
其中几个直接读运行态与执行历史:
| 模板 | 解析成 | 源码 |
|---|---|---|
{{ $flow.state.* }} | flowConfig.state(即 agentflowRuntime.state) | :390 |
{{ $form.* }} | 当前/运行态表单值 | :303 |
{{ $webhook.body.* }} 等 | webhook 载荷 | :315 |
{{ $iteration.* }} | 当前迭代项(迭代子流内) | :356 |
{{ $loopCount }} | 最近 loop 节点当前圈数(读 loopCounts) | :338 |
{{ nodeId.output.path }} | 从 agentFlowExecutedData 里翻该节点最近一次输出 | :402 |
注意最后一条:引用上游节点的输出,是去执行历史里"倒着找最近一次匹配"(:411、:438),
而不是顺着边拿——这正是解释执行模型下"节点间不直连数据"的取数方式。
6.3 webhook 引用的预解析
webhook 触发时,起点节点的 webhookDefaultInput 模板里的 {{ $webhook.* }} 在任何节点跑之前
就先被 resolveWebhookRefs(buildAgentflow.ts:110)解析一遍(:1763),保证起点 run() 和下游
看到的是同一个值。它还防原型链穿越(:118 挡掉 __proto__/constructor/prototype)。
7. 特色能力(各节点如何借引擎实现控制流)
引擎本身只提供"队列 + 等待室 + 计数器 + 剪枝"这几样原语。真正的控制流语义,是各节点在 run() 里
返回特定形状的输出、再由引擎解读出来的。下面这些节点都在
packages/components/nodes/agentflow/*。
7.1 Condition / ConditionAgent —— 条件分支
- Condition(
Condition/Condition.ts:274):按用户配的 if/else 规则逐条判定,命中的置isFulfilled:true;一条都没中就补一个 else 分支置真(:340-359)。输出output.conditions正是 §5.3 剪枝的依据。 - ConditionAgent(
ConditionAgent/ConditionAgent.ts:249):把"走哪条"交给 LLM 判断, 输出同样是一组带isFulfilled的 conditions,下游剪枝逻辑复用同一套。
7.2 Loop / Iteration —— 两种循环
| 节点 | 循环方式 | 关键输出 |
|---|---|---|
Loop(Loop/Loop.ts:104) | 回跳式:输出 nodeID(跳回哪个节点)+ maxLoopCount,引擎把该节点重新入队 | output.nodeID(:143) |
Iteration(Iteration/Iteration.ts:39) | 遍历式:输出一个数组,引擎对每个元素递归跑子流 | output.iterationInput |
Iteration 的递归执行在 executeNode 里(buildAgentflow.ts:1262-1406):为迭代块内的子节点
单独造一份 flowData,逐元素带 iterationContext 调 executeAgentFlow(..., isRecursive:true),
再把子流结果和运行态合并回父流。
7.3 HumanInput —— 人在环暂停/恢复
HumanInput/HumanInput.ts:136。一个节点,两种角色,看有没有带 humanInput:
第一次跑到(没带 humanInput): 人点了按钮再调一次(带 humanInput):
├ 流式发出问题描述 ├ run() 走 humanInput 分支 (HumanInput.ts:155)
├ engine 造 "通过/拒绝" action ├ 按 type 把 proceed/reject 之一置 isFulfilled
├ 标 STOPPED, 存进 Execution 表 ├ 引擎从上次 STOPPED 存档恢复 (1775-1873)
└ return shouldStop:true → 主循环 break └ 从该节点续跑, 剪掉没选的分支
暂停侧在引擎的 executeNode(buildAgentflow.ts:1409-1452);恢复侧在 executeAgentFlow 开头一大段
"找上次 STOPPED 节点、校验、恢复 state、从它续跑"(:1775-1873)。这正是"长时运行"能成立的地方——
状态全在数据库里,跨请求。
7.4 ExecuteFlow —— 子流调用
ExecuteFlow/ExecuteFlow.ts:162。把另一个 agentflow当一步来调:向
/api/v1/prediction/{selectedFlowId} 发预测请求(:197),可选把子流结果并进本流 state
(:238)。它拦截"调用自己"防递归自噬(:189)。
7.5 DirectReply —— 直接回话
DirectReply/DirectReply.ts:40。最简单的节点:把配置好的文本直接流式发给前端
(streamTokenEvent,:50)并作为 output.content 返回——用于"不经 LLM、直接回一句固定话"。
8. 执行持久化(回看与恢复的地基)
V2 把每一步都写进数据库的 Execution 表。这既支撑前端"执行回放",也是人在环恢复的存档。
- 创建:非递归执行一开始就
addExecution(buildAgentflow.ts:184),插一条state=INPROGRESS、executionData为 JSON 化的agentFlowExecutedData的记录(:1899)。 - 更新:出错、迭代进展、人在环暂停、整体收尾,都调
updateExecution(buildAgentflow.ts:212) 覆盖executionData和state;置STOPPED时还盖上stoppedDate(:230-232)。收尾处按 执行数据里有没有 TERMINATED/ERROR/STOPPED 决定最终state(:2213-2233)。 agentFlowExecutedData是什么:一个数组,每个节点跑完 push 一项,含nodeId、nodeLabel、 节点结果data、previousNodeIds、status(:2110)。它既是落库内容,也是 §6.2 里{{ nodeId.output }}取数的来源——执行历史本身就是数据总线。
递归子流(迭代块)不新建执行记录,而是复用 parentExecutionId 把数据并回父记录(:1360-1369、
:1874-1892)。
9. 边界与坑(诚实)
- 单起点:非递归执行只允许一个起点节点,多起点直接抛错(
buildAgentflow.ts:1519)。 - 两道循环护栏:
MAX_LOOP_COUNT(默认 10,管 loop 回跳,:174)与MAX_ITERATIONS(默认 1000,管主循环总轮数,:1912)。设计上宁可提前停也不放任跑飞。 - 恢复态严格:只有上次是
STOPPED(或STOPPED后紧跟一个ERROR)的执行才能续跑, 其它状态直接拒绝(:1783-1818)——防止从一个不一致的中间态胡乱恢复。 - 恢复要求图没改:续跑时校验
startNodeId仍存在于上次执行数据,否则报"可能是无效恢复或图被改过" (:1836-1843)。 - 条件全不命中会强走最后一条:这是刻意的兜底(§5.3),但也意味着"最后一条分支"隐含承担了 else 语义,配图时要留意。
10. 代码地图(导航索引)
引擎主体(packages/server/src/utils/buildAgentflow.ts):
| 主题 | 符号名 | 行 |
|---|---|---|
| 引擎入口 / 主循环 | executeAgentFlow | 1541 / 循环 2020 |
| 待办队列元素类型 | INodeQueue | 75 |
| 等待室元素类型 | IWaitingNode | 67 |
| 三容器初始化 | nodeExecutionQueue/waitingNodes/loopCounts | 1650-1652 |
| 循环上限常量 | MAX_LOOP_COUNT | 174 |
| 单节点执行 | executeNode | 1051 |
| 分流 / 下一步决策 | processNodeOutputs | 847 |
| 条件剪枝 | determineNodesToIgnore | 799 |
| 就绪判定 | hasReceivedRequiredInputs | 770 |
| 依赖分析 | setupNodeDependencies | 672 |
| 多入边输入合并 | combineNodeInputs | 971 |
| 运行态类型 | IAgentFlowRuntime | 97 |
| 变量解析 | resolveVariables | 241 |
| webhook 引用预解析 | resolveWebhookRefs | 110 |
| 执行落库 | addExecution / updateExecution | 184 / 212 |
特色节点(packages/components/nodes/agentflow/):
| 节点 | 符号 / 文件 | run 行 |
|---|---|---|
| 条件分支 | Condition/Condition.ts (conditionAgentflow) | 274 |
| LLM 条件分支 | ConditionAgent/ConditionAgent.ts (conditionAgentAgentflow) | 249 |
| 回跳循环 | Loop/Loop.ts (loopAgentflow) | 104 |
| 遍历 循环 | Iteration/Iteration.ts (iterationAgentflow) | 39 |
| 人在环 | HumanInput/HumanInput.ts (humanInputAgentflow) | 136 |
| 子流调用 | ExecuteFlow/ExecuteFlow.ts (executeFlowAgentflow) | 162 |
| 直接回话 | DirectReply/DirectReply.ts (directReplyAgentflow) | 40 |
同组其它章: Flowise 全景 · 数据模型与节点系统 · 经典执行引擎 · 请求生命周期与流式/队列 · 前端画布。节点接口通用形态见 01,SSE 流式与队列扩展见 04,本章不重复。