数据截至 (上游 commit 7a975c596eca)
环、循环与条件路由 — 这张图为什么不是 DAG
30 秒导读: 前一章(调度引擎)讲的是"按拓扑序把顶点一个个跑完"。这章讲 Langflow 在这套 DAG 调度之上多加的两件事——允许图里有环(节点可以回跳到上游重跑)和允许运行期剪枝 (一个 If-Else 判完之后,没选中的那半张图整个不跑)。这两件事各自都会打破"所有前驱跑完才能跑我" 这条 DAG 铁律,所以引擎里为它们准备了单独的判据和单独的状态机。
1. 这节讲什么:普通 DAG 引擎缺的两样东西
一个普通的 DAG(有向无环图)工作流引擎有两条硬规矩:
- 不许有环——有环就没法拓扑排序,调度器不知道谁先谁后。
- 顶点全跑——图上画了的节点,只要前驱跑完了就得跑。
画布上的用户偏偏想要这两件被禁止的事:
| 用户想要的 | 画布上长什么样 | 打破了哪条规矩 |
|---|---|---|
| "答案不满意就回去重问一 次" | 下游节点连回上游节点,形成回边 | 不许有环 |
| "条件为真走这半边,为假走那半边" | 一个 If-Else 分出两条支路,只该跑一条 | 顶点全跑 |
| "对列表里每一项都跑一遍这段流程" | Loop 的 item 输出接一段子流程再接回来 | 两条都破 |
Langflow 的做法不是"把环拆掉再当 DAG 跑",而是保留环、给环上的顶点一套单独的放行规则; 不是"跑完再丢弃没用的结果",而是在调度阶段就把不该跑的分支标记成不可运行。
一句话直觉: 把 DAG 引擎想成"红绿灯按拓扑序依次放行"。Langflow 干了两件事——给环上的路口 换了一套放行逻辑(否则环上每个路口都在等对面先走,谁也走不了),再给分支路口装了可以临时封路的路障。
2. 顶层全景:一次带环带分支的运行
先看一张最小的、同时有环和有分支的流:一个 If-Else 判断结果够不够好,不够好就把消息拼一拼再回去重判, 够好了才往输出走。
怎么读这张图: 实线是普通数据流,↰ 那条是回边(构成环);虚线框里是被剪掉的分支。从左往右是主流向。
┌──────────────┐
ChatInput ───────▶ │ If-Else │ ──true──▶ TextOutput ──▶ ChatOutput
▲ │ (路由顶点) │
│ └──────┬───────┘
│ │ false
│ ▼
│ ┌─────────────┐
└──── ↰ ───────│ Concatenate │ ← 这条回边让 ChatInput / Concatenate /
回边 └─────────────┘ If-Else 三个顶点同处一个环
每一轮判断结束时,引擎要同时做两件事:
- 剪掉没选中的那半边(true 走了就把 false 那支封掉,反之亦然);
- 决定环要不要再转一圈(false 分支通了,就意味着要回到 ChatInput 重跑)。
这两件事在代码里由两套完全独立的机制完成。下面这张表是本章的骨架,后面每一节展开其中一 格:
| 关切 | 引擎里的机制 | 核心符号 | 生命周期 |
|---|---|---|---|
| 谁在环上 | 强连通分量检测 | find_cycle_vertices | 建图时算一次并缓存 |
| 环上顶点何时能跑 | 两套前驱判据 | are_all_predecessors_fulfilled | 每次询问时现算 |
| 环怎么停下来 | 每顶点产出计数 + 上限 | should_continue / yielded_counts | 整个 run 期间累加 |
| 分支封路(服务于环) | ACTIVE / INACTIVE 顶点状态 | mark_branch | 每步结束即重置 |
| 分支封路(服务于路由) | 条件排除集合 | conditionally_excluded_vertices | 由源顶点持有,直到它重新判定 |
最容易踩的坑就在最后两行:Langflow 里并存着两套剪枝状态机,它们看起来在做同一件事, 但生命周期完全相反。第 5 节专门讲这个。
3. 环:怎么识别、怎么标记
3.1 谁在环上 —— 强连通分量
判断"图里有没有环"和"具体哪些顶点在环上"是两个问题,Langflow 的工具箱里四个函数各 管一摊:
| 函数 | 回答什么问题 | 算法 | 位置 |
|---|---|---|---|
has_cycle | 有没有环(布尔) | DFS + 递归栈 | utils.py:408 |
find_cycle_edge | 第一条造成环的边 | DFS,命中即返回 | utils.py:444 |
find_all_cycle_edges | 所有造成环的边 | DFS,全部收集 | utils.py:481 |
find_cycle_vertices | 哪些顶点在环上 | networkx 强连通分量 | utils.py:524 |
真正被运行期依赖的是最后一个。它不用 DFS,而是借 networkx 求强连通分量(SCC,一组互相都能到达的
顶点),分量里顶点多于一个、或者有自环,就整组算作环上顶点
(src/lfx/src/lfx/graph/graph/utils.py:531-533,find_cycle_vertices):
for component in nx.strongly_connected_components(graph):
if len(component) > 1 or graph.has_edge(tuple(component)[0], tuple(component)[0]):
cycle_vertices.update(component)
为什么用 SCC 而不是 DFS? DFS 找到的是"回边"(哪条边造成了环),但运行期真正需要的是"这个顶点是不是 在环里"——因为放行判据是按顶点问的。SCC 天然给出的就是顶点集合,而且不依赖从哪个入口开始遍历。
Graph 上三个属性把这些包成缓存:
| 属性 | 内容 | 位置 |
|---|---|---|
Graph.cycle_vertices | 环上顶点集合(懒算 + 缓存) | base.py:2204-2209 |
Graph.is_cyclic | 就是 bool(self.cycle_vertices) | base.py:645-653 |
Graph.cycles | 造成环的边列表,需要 _start 才算 | base.py:2193-2202 |
注意 cycles 和 cycle_vertices 用的不是同一套算法:前者调 find_all_cycle_edges,必须有起点顶点
self._start,没有就返回空列表;后者不需要起点。所以从 JSON 载入的流(没有显式 start/end)上
graph.cycles 恒为空,但 graph.cycle_vertices 照常工作——运行期依赖的是后者。
缓存会在 add_nodes_and_edges 里被两次清空(base.py:276-277、base.py:282-284),因为图结构变了
环也就变了。
3.2 环上的边换一个类:CycleEdge
建边的时候,只要两端有任意一端落在 cycle_vertices 里,这条边就不是普通 Edge 而是 CycleEdge
(src/lfx/src/lfx/graph/graph/base.py:2612-2615,Graph.build_edge):
if any(v in self.cycle_vertices for v in [source.id, target.id]):
new_edge: CycleEdge | Edge = CycleEdge(source, target, edge)
else:
new_edge = Edge(source, target, edge)
两个类的差别很小但很关键:
Edge | CycleEdge | |
|---|---|---|
is_cycle | False(edge/base.py:52) | True(edge/base.py:291) |
| 给两端顶点打标 | 无 | source.has_cycle_edges = True,target 同样(edge/base.py:292-293) |
| 额外能力 | 无 | honor() / get_result_from_source()(edge/base.py:295-332) |
CycleEdge.honor 做的事是"履约":把已经建好的源顶点的结果塞进目标顶点的 params,并置 is_fulfilled
(src/lfx/src/lfx/graph/edge/base.py:295-318)。它明确拒绝在这里触发构建——源没 built 就直接抛错,
注释写得很直白:这条路径必须是只读的。
is_cycle 这个标志还有一个下游用处:当某个前驱始终没建成时,ComponentVertex._get_result 会沿着
cycle 边去取模板默认值,而不是报错(见第 6 节)。
3.3 环上的顶点强制关缓存
环的意义就是"同一个顶点跑第二遍要得到新结果",所以建图时会把环上所有顶点的输出缓存关掉
(src/lfx/src/lfx/graph/graph/base.py:1886-1892,Graph._set_cache_to_vertices_in_cycle):
cycle_vertices = set(find_cycle_vertices(edges))
for vertex in self.vertices:
if vertex.id in cycle_vertices:
vertex.apply_on_outputs(lambda output_object: setattr(output_object, "cache", False))
它在 _build_graph 里被调用(base.py:1487),紧接着的循环把环上顶点登记进
run_manager.cycle_vertices(base.py:1489-1491),prepare() 里还会再登记一次首层里的环上顶点
(base.py:2301-2304)。调度器判断"这个顶点在不在环上"读的是 run_manager.cycle_vertices 这份副本,
不是 Graph.cycle_vertices——两者由上面这些调用保持同步。
顺带一提,冻结(frozen)机制也给环开了口子:Loop 类顶点即使被冻结也必须重建
(base.py:1717-1718,is_loop_component = vertex.display_name == "Loop" or vertex.is_loop)。
3.4 拓扑排序怎么给环找一个入口
环上每个顶点的入度都 ≥ 1,标准的 Kahn 算法开局就找不到"入度为 0"的顶点,队列是空的。
layered_topological_sort 对这种情况有专门分支(src/lfx/src/lfx/graph/graph/utils.py:568-585):
is_cyclic 且所有顶点入度 > 0
│
├── 有 start_id ──────────▶ 队列 = [start_id],并把它的入度强行改成 0
│
└── 没有 start_id ────────▶ find_start_component_id 找 webhook/chat 类输入顶点
找不到就随便挑一个顶点当入口
排序阶段还允许环上顶点在层里重复出现最多两次——常量 MAX_CYCLE_APPEARANCES = 2
(utils.py:11),判据在 utils.py:645 和 utils.py:670-673。这只是给运行期铺个路,真正决定跑多少圈的
是运行期的判据和止损计数,不是这份静态分层。分层本身的细节见 调度引擎。
3.5 止损:环靠什么停下来
Langflow 不做静态的循环次数分析,它靠一个朴素的运行期计数器兜底。
async_start 每 yield 一个结果就给该顶点的计数加一(src/lfx/src/lfx/graph/graph/base.py:496-506):
yielded_counts: dict[str, int] = defaultdict(int)
while should_continue(yielded_counts, max_iterations):
result = await self.astep(...)
yield result
if isinstance(result, Finish):
return
if hasattr(result, "vertex"):
yielded_counts[result.vertex.id] += 1
should_continue 的判据是"产出次数最多的那个顶点还没超过上限"
(src/lfx/src/lfx/graph/graph/utils.py:518-521):
def should_continue(yielded_counts: dict[str, int], max_iterations: int | None) -> bool:
if max_iterations is None:
return True
return max(yielded_counts.values(), default=0) <= max_iterations
三个容易搞错的细节:
-
正常结束走的是
return,不是循环条件。 队列空了astep返回Finish,生成器直接 return (base.py:401-402)。只有真的转超了才会走到循环外面那句raise ValueError("Max iterations reached")(base.py:420-421)。 -
计数是"每顶点"的,不是"总步数"。 一条长直线流跑 50 个不同顶点,每个计数都是 1,不会触发上限。
-
同步
start()强制要求max_iterations,异步async_start()不要求。 只有Graph.start里有这道前置检查(base.py:464-466):if self.is_cyclic and max_iterations is None:msg = "You must specify a max_iterations if the graph is cyclic"raise ValueError(msg)也就是说,直接调
async_start且不传max_iterations的环流,理论上会无限转——should_continue在max_iterations is None时恒为True。真实的兜底通常来自组件层自己的max_iterations(见 7.1 的 If-Else)。
回归测试直接锁住了这个行为:src/lfx/tests/unit/graph/graph/test_cycles.py:114-115 用
max_iterations=2 跑一个必然转不完的环,断言抛出 Max iterations reached。