数据截至 (上游 commit 7a975c596eca)
03 · 调度引擎:谁先跑、谁能并行、结果怎么往下传
本章讲什么: 一张图被构造成对象之后,Langflow 靠什么决定「先跑哪个组件、哪些可以同时跑、跑完的值怎么到下一个组件手里」。四个环节:排序 → 就绪判定 → 并发构建 → 结果拉取。
路径约定(重要): 本章所有以
graph/开头的路径,完整前缀都是src/lfx/src/lfx/。 例如graph/graph/base.py:2664的真实路径是src/lfx/src/lfx/graph/graph/base.py:2664。 「代码地图」一节给的是完整路径。
3.1 先想清楚:调度器到底要回答什么
画布上的一张流,本质是一堆组件(顶点)加连线(边)。点一下 Run,系统必须回答三个问题:
| 问题 | 白话 | 本章对应小节 |
|---|---|---|
| 谁先跑? | 没有输入依赖的组件先跑,有依赖的等前面的 | §3.3 排序 |
| 谁现在能跑? | 一个组件的所有上游都出结果了吗? | §3.4 就绪判定 |
| 谁能同时跑? | 互不依赖的组件为什么要排队 | §3.5 并发构建 |
| 值怎么过去? | Prompt 组件怎么拿到 ChatInput 的那句话 | §3.6 结果拉取 |
一句话直觉: 把它想成厨房出菜。排序 = 排出「洗菜 → 切菜 → 下锅」的大致班次;就绪判定 = 每道菜自己数「我要的料到齐了没」;并发构建 = 同一班次的几口锅一起开火;结果拉取 = 下一道菜的厨师自己去上一个灶台端盘子,而不是上一个厨师送过来。
最后那句是 Langflow 最容易被误解的一点,§3.6 会展开。
3.2 顶层全景:四步流水线
这张图从上到下读。左边一列是「一次性的准备」,右边的循环是「反复推进」。
Graph.prepare()
│
┌───────────────┴───────────────┐
│ ① 排序 │ 出两样东西:
│ 分层拓扑 + 层内再排序 │ · first_layer(起跑名单)
│ get_sorted_vertices │ · vertices_to_run(允许跑的全集)
└───────────────┬───────────────┘
│ first_layer 灌进 _run_queue
▼
┌───────────────────────────────────────────────┐
│ 执行循环(反复) │
│ │
│ ② 就绪判定 ──► ③ 构建 ──► ④ 结果拉取 │
│ 四个集合对账 gather 并发 下游 await 上游 │
│ ▲ │ │
│ └──── 新就绪顶点回填队列 ◄────┘ │
└───────────────────────────────────────────────┘
两条执行路径,同一套内核
Langflow 有两个入口驱动这个循环,内核完全共用,区别只在「一次推进多少个顶点」:
| 路径 | 入口 | 一次推进 | 并发 | 典型用途 |
|---|---|---|---|---|
| 单步流 | Graph.astep(graph/graph/base.py:1949) | 1 个顶点 | 无(串行) | async_start 逐步 yield,前端逐节点渲染 |
| 批处理 | Graph.process(graph/graph/base.py:2187) | 一整层 | asyncio.gather | arun / _run 一次性跑完整张流 |
两者调的都是同一个 build_vertex 和同一个 get_next_runnable_vertices。
部件一句话职责
| 部件 | 干什么 | 在哪 |
|---|---|---|
get_sorted_vertices | 剪枝 + 分层 + 层内排序,产出 (first_layer, remaining_layers) | graph/graph/utils.py:870 |
Graph.prepare | 调排序、把首层灌进运行队列、打第一个快照 | graph/graph/base.py:2664 |
RunnableVerticesManager | 用四个集合记「谁欠谁」,回答「这个顶点现在能跑吗」 | graph/graph/runnable_vertices_manager.py:4 |
Graph.build_vertex | 冻结缓存判断 + 调 Vertex.build + 回写缓存 | graph/graph/base.py:2049 |
Graph._execute_tasks | 一层任务并发跑,收结果、算下一批 | graph/graph/base.py:2342 |
Vertex._build_each_vertex_in_params_dict | 构建前把参数里的上游顶点换成上游的真实结果 | graph/vertex/base.py:554 |
ComponentVertex._get_result | 按边的 handle 名,从上游众多 output 里挑出对的那个 | graph/vertex/vertex_types.py:93 |
3.3 第一步:排序 —— 排出班次
3.3.1 它要解决的小问题
拓扑排序(topological sort,按依赖把有向图排成线性顺序)只给你一条串行的链。但画布上经常有三个互不相干的分支,它们本可以同时跑。所以 Langflow 用的是分层拓扑排序:同一层的顶点彼此无依赖,可以并发;层与层之间才有先后 。
入口是一个函数:get_sorted_vertices(graph/graph/utils.py:870),返回 (first_layer, remaining_layers)。
3.3.2 排序前先剪枝:只跑用得着的那部分
用户可能只想「跑到某个组件为止」(stop)或者「从某个组件开始跑」(start)。这靠两个对称的 BFS 完成:
| 函数 | 方向 | 语义 | 位置 |
|---|---|---|---|
filter_vertices_up_to_vertex | 沿前驱回溯 | 目标顶点的所有祖先(算它需要谁) | graph/graph/utils.py:1010 |
filter_vertices_from_vertex | 沿后继前进 | 起点顶点的所有后代(它能影响谁) | graph/graph/utils.py:1066 |
两者都是 deque + visited 的标准广度优先,唯一的讲究是只收 vertices_ids 里的顶点——已经被上一轮过滤掉的不再拉回来。
start_component_id 的处理多绕一层(graph/graph/utils.py:951-964):先求出「从 start 可达的顶点集」,再对每一个可达顶点求它的祖先集并起来。原因很实在——从 start 出发能走到 LLM,但 LLM 还需要一个 Prompt 模板顶点作为输入,而那个 Prompt 从 start 走不到。只取后代会漏掉它。
还有一个专治环的小动作:如果 stop 组件本身在环上,就把它改当作 start(graph/graph/utils.py:905-907)。环的完整 讨论见第 04 章。
3.3.3 分层:入度归零就进下一层
layered_topological_sort(graph/graph/utils.py:538)是经典的 Kahn 算法加了几处补丁。先看主干思路:
# 示意,非源码 —— 演示分层拓扑的骨架
queue = [v for v in vertices if in_degree[v] == 0] # 无人依赖的先入队
while queue:
layer = queue # 当前这批就是一层
queue = []
for v in layer:
for nxt in successors[v]:
in_degree[nxt] -= 1 # 拆掉一条入边
if in_degree[nxt] == 0:
queue.append(nxt) # 入度归零 → 下一层
layers.append(layer)
真源码在这个骨架上加了四处补丁,每一处都对应一类真实图形:
| 补丁 | 干什么 | 位置 |
|---|---|---|
| 全员入度 > 0 的兜底 | 整张图都在环里时没有天然起点,就强行把 start_id(或找到的 ChatInput)入度置 0 当起点 | graph/graph/utils.py:568-586 |
| 首层单独处理 | 用 first_layer_vertices 集合去重,防同一顶点在第一层出现两次 | graph/graph/utils.py:601-638 |
| 补捞前驱 | 邻居入度减完仍 > 0 时,反过来把它那些已就绪或在环上的前驱也塞进队列,避免整队卡死 | graph/graph/utils.py:621-630、668-674 |
| 环顶点限次复现 | 环上顶点允许重复出现在多层,上限由 MAX_CYCLE_APPEARANCES = 2 卡死(graph/graph/utils.py:11) | graph/graph/utils.py:645 |
找起点用的 find_start_component_id(graph/graph/utils.py:14)——它优先挑 webhook,其次挑 ChatInput。
3.3.4 层内再排序:两条规则
分完层还要决定层内顺序。两个函数依次作用(graph/graph/utils.py:1000-1002):
规则一:ChatInput 单独站第一层。 sort_chat_inputs_first(graph/graph/utils.py:815)按 vertex id 里是否含 "ChatInput" 判断,把它抠出来自成一层放最前。三个细节值得记:
- 如果 ChatInput 有前驱(说明它被别人喂数据,不是真起点),直接原样返回不动(
graph/graph/utils.py:839-840)。 - 全图只允许一个 ChatInput,第二个直接抛
ValueError(graph/graph/utils.py:843-845)。 - 只有在没指定
start_component_id时才做这步(graph/graph/utils.py:999)——用户手动指了起点,就不再越权重排。
规则二:层内按「依赖深度」降序。 sort_layer_by_dependency(graph/graph/utils.py:799)对每层调 _sort_single_layer_by_dependency(graph/graph/utils.py:751),核心是递归算「本顶点的后继在本层中的最大下标」,然后 sorted(..., reverse=True)(graph/graph/utils.py:796)。
这里有个防死循环的巧劲:递归时用 processing 集合记录「正在算的顶点」,一旦重入就说明这一层内部有环,直接返回当前下标把环切断(graph/graph/utils.py:781-786)。注释里写得很明白——层内出现环是环图的正常现象,不是前面环检测的 bug。
3.3.5 落到 Graph 上:sort_vertices 与 prepare
Graph.sort_vertices(graph/graph/base.py:2748)是薄封装:先把所有 顶点标回 ACTIVE(base.py:2376),把图的四张邻接表喂给 get_sorted_vertices(base.py:2378),然后把结果摊平:
# 真实源码节选 graph/graph/base.py:2394-2398 —— 重点看 vertices_to_run 是一个"全集"
self.vertices_layers = remaining_layers
self.vertices_to_run = set(chain.from_iterable([first_layer, *remaining_layers]))
self.build_run_map()
self._first_layer = first_layer
Graph.prepare(graph/graph/base.py:2664)在这之上做四件事:把首层每个顶点登记进 vertices_being_run(顺带把环上顶点登记进 cycle_vertices,base.py:2301-2303)、_first_layer = sorted(first_layer)、首层灌进 _run_queue、打第一张快照(base.py:2305-2308)。
排序失败还有兜底:指定了 start/stop 但排序抛异常时,prepare 会退回无参数的全图排序(graph/graph/base.py:2671-2674)。
3.3.6 一个反直觉的事实:分层结果大部分被丢掉了
读到这里很容易以为「执行就是按 layers 一层层跑」。不是。 顺着 self.vertices_layers(即 remaining_layers)的消费点查一遍,它只出现在三个地方:
| 消费点 | 用途 | 位置 |
|---|---|---|
next_vertex_to_build | 一个 generator,全仓库无任何调用者 | graph/graph/base.py:1161 |
get_snapshot / _snapshot | 调试快照 | graph/graph/base.py:2018、427 |
__getstate__ | 序列化 | graph/graph/base.py:1410 |
也就是说,真正进入执行决策的只有两样:first_layer(起跑名单)和 vertices_to_run(允许跑的全集)。之后每一步「下一个跑谁」,全部由 §3.4 的前驱对账现算。分层排序的实际职责是「挑起跑线 + 圈定范围」,而不是「排好执行时刻表」。
顺带一提,utils.py 和 base.py 里还躺着三个没有任何调用者的排序函数——refine_layers(graph/graph/utils.py:682)、Graph.sort_interface_components_first(graph/graph/base.py:2780)、Graph.sort_by_avg_build_time(graph/graph/base.py:2795)。最后那个想按历史平均构建耗时排序(数据源是 Vertex.avg_build_time,graph/vertex/base.py:162),思路不错但没接上。
3.4 第二步:就绪判定 —— 四个集合的对账
3.4.1 它要解决的小问题
「A 跑完了,接下来谁能跑?」这个问题不能靠查图现算——因为得考虑「B 的另一个上游 C 还没跑完」「B 已经被别的任务领走了」「B 根本不在本次要跑的范围里」。所以 Langflow 用一个专门的小对象记账。
整个文件只有 133 行:graph/graph/runnable_vertices_manager.py。
3.4.2 四个集合各管一件事
| 集合 | 类型 | 含义 | 定义处 |
|---|---|---|---|
run_map | dict[str, list[str]] | 谁的后继是谁(前驱 → 后继列表) | runnable_vertices_manager.py:6 |
run_predecessors | dict[str, list[str]] | 每个顶点还欠着哪些前驱,跑完一个删一个 | runnable_vertices_manager.py:7 |
vertices_to_run | set[str] | 本次运行允许跑的全集(来自排序阶段圈定的范围) | runnable_vertices_manager.py:8 |
vertices_being_run | set[str] | 已被领走(正在跑或已排入队列)的顶点,防重复调度 | runnable_vertices_manager.py:9 |
另有两个服务于环的集合:cycle_vertices 和 ran_at_least_once(runnable_vertices_manager.py:10-11),用法归第 04 章。
run_map 是从 run_predecessors 反向翻出来的,一次建好:
# 真实源码节选 graph/graph/runnable_vertices_manager.py:111-116 —— 前驱表反转成后继表
self.run_map = defaultdict(list)
for vertex_id, predecessors in predecessor_map.items():
for predecessor in predecessors:
self.run_map[predecessor].append(vertex_id)
self.run_predecessors = predecessor_map.copy()
self.vertices_to_run = vertices_to_run
Graph.build_run_map(graph/graph/base.py:2817)就是把图的 predecessor_map 和 vertices_to_run 递进来触发这个反转。
3.4.3 三道闸 + 一次清算
is_vertex_runnable(runnable_vertices_manager.py:56)的判断顺序很直白:
vertex_id
│
┌────▼────┐ 否
│ 是 ACTIVE?├────────► 不能跑
└────┬────┘
是 │
┌────▼──────────┐ 是
│ 已被别人领走? ├────────► 不能跑
│ vertices_being_run
└────┬──────────┘
否 │
┌────▼──────────────┐ 否
│ 在允许跑的全集里? ├────────► 不能跑
│ vertices_to_run │
└────┬──────────────┘
是 │
┌────▼──────────────────┐
│ run_predecessors 空了? │──► 空 = 能跑
└───────────────────────┘ 非空 = 环规则另判(第 04 章)
最后一道「前驱清空了没」在 are_all_predecessors_fulfilled(runnable_vertices_manager.py:67)。非环顶点的规则只有一句:pending 为空才行,否则 return False(runnable_vertices_manager.py:100)。
「清算」发生在顶点跑完时——remove_from_predecessors(runnable_vertices_manager.py:102)顺着 run_map 找到自己所有后继,把自己从每个后继的欠 账表里划掉:
# 真实源码节选 graph/graph/runnable_vertices_manager.py:104-107 —— 跑完就销账
predecessors = self.run_map.get(vertex_id, [])
for predecessor in predecessors:
if vertex_id in self.run_predecessors[predecessor]:
self.run_predecessors[predecessor].remove(vertex_id)
(变量名 predecessors 其实装的是后继,读的时候别被绕进去。)
remove_vertex_from_runnables(runnable_vertices_manager.py:125)把「从 vertices_being_run 移除」和「销账」打包成一个动作。
Graph.is_vertex_runnable(graph/graph/base.py:2808)在转发给 manager 之前,还先查一道 conditionally_excluded_vertices——条件路由剪掉的分支不许跑,细节见第 04 章。
3.4.4 找下一批:不只看后继,还会回头捞前驱
直觉做法是「跑完 A,看 A 的后继谁就绪了」。Langflow 多做了一步。
find_next_runnable_vertices(graph/graph/base.py:2253)遍历 A 的后继:后继就绪 → 收进来;后继没就绪 → 调 find_runnable_predecessors_for_successor(graph/graph/base.py:2837)去递归回溯这个后继的前驱链,把链上任何已就绪的顶点捞出来。
为什么需要这一步?看这张图:
A ──────────────► D
▲
B ───► C ─────────┘
(B 与 A 同在首层,C 不在首层)
A 跑完,D 还欠着 C,不就绪。若只看后继就到此为止,而 C 恰好因为不在首层从没被调度过——整张流卡死。回溯逻辑此时沿 D → C 发现 C 已就绪(B 早跑完了),把 C 捞出来接着跑。
visited 集合防止回溯在环里打转(graph/graph/base.py:2839、2464-2466)。
3.4.5 提交:get_next_runnable_vertices
顶点跑完后的收尾统一走 get_next_runnable_vertices(graph/graph/base.py:2268),四个动作在同一把锁里完成:
- 把自己从 runnables 里摘掉并销账(
base.py:1916); - 算出下一批(
base.py:1917); - 把下一批逐个塞进
vertices_being_run预占(base.py:1919-1923)——这是并发下防重复调度的关键,预占后别的任务再算is_vertex_runnable就会在第二道闸被拦下; - 可选地把整张 Graph 写回缓存(
base.py:1924-1926)。
自己算出自己的情况会被剔除(base.py:1920-1921);状态顶点(state vertex)还会额外把 activated_vertices 追加进来(base.py:1927-1928)。
3.5 第三步:执行 —— 两条路,一个内核
3.5.1 单步路径:astep + 运行队列
_run_queue 是一个 deque(graph/graph/base.py:181),prepare 把首层灌进去。之后每次 astep(graph/graph/base.py:1949)干六件事:
| 步骤 | 做什么 | 位置 |
|---|---|---|
| 1 | 队列空 → 结束追踪,返回 Finish() | base.py:1592-1594 |
| 2 | get_next_in_queue() 从队头取一个 | base.py:1572、1595 |
| 3 | build_vertex(...) 真正构建 | base.py:1618 |
| 4 | get_next_runnable_vertices(...) 算下一批 | base.py:1629 |
| 5 | extend_run_queue(...) 回填队尾 | base.py:1577、1634 |
| 6 | 重置临时失活/激活标记,写缓存,打快照 | base.py:1635-1640 |
一个小特权:如果本次设了 stop_vertex 且它出现在下一批里,下一批就只剩它(base.py:1632-1633)——到此为止,别的分支不再跑。
async_start(graph/graph/base.py:419)把 astep 包成异步生成器:prepare → initialize_run → 循环 astep 并 yield 每步结果,遇到 Finish 收工(base.py:396-402)。防跑飞的闸是 should_continue(graph/graph/utils.py:518),按「同一顶点被 yield 的最大次数」和 max_iterations 比较,超了抛 "Max iterations reached"(base.py:420-421)。
3.5.2 批处理路径:process 一层一层并发
Graph.process(graph/graph/base.py:2187)是另一种驱动方式 ——把「当前这批」整个铺开并发跑:
# 示意,非源码 —— process 的主循环骨架
to_process = deque(first_layer)
while to_process:
batch, tasks = list(to_process), []
to_process.clear()
for vertex_id in batch:
tasks.append(asyncio.create_task(self.build_vertex(vertex_id=vertex_id, ...)))
next_ids = await self._execute_tasks(tasks, lock=lock) # 等这一层全部结束
if not next_ids:
break # 没有新就绪的 → 收工
to_process.extend(next_ids)
对照真源码:current_batch 复制与清空在 base.py:1847-1848,建任务在 base.py:1852-1863,_execute_tasks 调用在 base.py:1869,空则 break 在 base.py:1875-1876。
任务名是当作数据通道用的。 建任务时名字写成 f"{vertex.id} Run {n}"(base.py:1862),出异常时再 get_name().split(" ")[0] 反解出 vertex id(base.py:1992)。能成立是因为 Langflow 的 vertex id 形如 ChatInput-Ab12,不含空格。
3.5.3 一层并发的收尾:_execute_tasks
_execute_tasks(graph/graph/base.py:2342)是并发的核心,读的时候盯住三件事。
其一,gather 收全再逐个查验。
# 真实源码节选 graph/graph/base.py:1985 —— 注意 return_exceptions=True
completed_tasks = await asyncio.gather(*tasks, return_exceptions=True)
return_exceptions=True 意味着异常不会当场炸出来,而是变成结果列表里的一个元素。随后循环里遇到 Exception 就取消剩余任务并原样 raise(base.py:2000-2002)。
这里有个值得注意的实现细节:asyncio.gather(..., return_exceptions=True) 会等所有任务都结束才返回,所以走到 for t in tasks[i + 1:]: t.cancel() 这行时,那些任务其实早就跑完了,cancel() 基本是空操作。换句话说「首个异常取消其余任务」在语义上成立(不会再往下推进),但不会真的把同层还在跑的组件掐断——同层的活儿一定会全部做完才报错。
其二,先全部销账,再统一算下一批。 两个循环是分开的:第一个循环把本层每个顶点从 runnables 里移除(base.py:2022-2028),第二个循环才逐个调 get_next_runnable_vertices(base.py:2030-2031)。注释说明了原因——并行顶点之间可能互为前驱/后继(典型是 ChatInput 这类输入顶点),先销完账再算,才不会把已经跑过的顶点又算成"就绪"重跑一遍。
其三,结果去重。 返回前 list(set(results))(base.py:2054),同一顶点被多个前驱同时点亮时只保留一份。
对外的成功日志与 SSE 事件也在第二个循环里发出(base.py:2034-2052),事件层细节见第 05 章。
_run(graph/graph/base.py:897)和 arun(graph/graph/base.py:1097)是这条路径的外壳:arun 负责把多组输入拆成多次运行,_run 负责灌输入、调 process、最后遍历所有已 built 的顶点收集输出(base.py:854-867)。
3.5.4 build_vertex:冻结组件直接复用上次的结果
Graph.build_vertex(graph/graph/base.py:2049)包在真正的 Vertex.build 外面,主要价值是一个缓存分支。
用户可以把一个组件标成 frozen(冻结)——意思是「这个组件的输出别再重算了,沿用上次的」。调试长流程时特别有用:改了最后一个组件,前面的 LLM 调用不必重跑重烧钱。
判定逻辑:
vertex.frozen 且不是 Loop 组件?
否 ─── ───────────► should_build = True(照常构建)
是
▼
get_cache(key=vertex.id)
miss ─────────────► should_build = True
hit
▼
把 built / artifacts / built_object / built_result / full_data / results
六个字段直接灌回顶点 ──► finalize_build() ──► 标记 used_frozen_result = True
(灌回时 KeyError 或 finalize 抛错 → 退回 should_build = True)
对应源码:frozen 判断 base.py:1718(注意 Loop 组件被排除,它必须每轮真跑),读缓存 base.py:1723,六字段回灌 base.py:1731-1736,used_frozen_result 标记 base.py:1742,两处失败回退 base.py:1744-1750。真构建后写回缓存在 base.py:1758-1769。
同样的 frozen 短路在顶点内部也有一份:Vertex.build 里 if self.frozen and self.built and not is_loop_component 直接返回既有结果(graph/vertex/base.py:835)。
最后包装成 VertexBuildResult(base.py:1785-1787),带上 result_dict、params、valid、artifacts、vertex 五项。没有结果则直接抛错(base.py:1780-1782)。
3.5.5 Vertex.build:顶点自己那一层
Vertex.build(graph/vertex/base.py:806)的执行主体只有一句 for step in self.steps(base.py:851-855),而 steps 初始化时就是 [self._build](graph/vertex/base.py:91)——这是一个留了扩展位但目前只有一步的设计。
整段跑在 async with self.lock 里(base.py:811),锁是每个顶点一把、懒初始化的(graph/vertex/base.py:120-124)。这把锁的作用是:同一顶点被多个下游同时拉取时,只会真构建一次。
进锁后的几道前置判断按顺序是:失活 → 返回 None(base.py:812-815);frozen 且已 built → 返回既有结果(base.py:820-821);已 built 且有 requester → 返回既有结果(base.py:822-825)。都不命中才 _reset() 重来。
_build(graph/vertex/base.py:393)的四步:拉上游结果(base.py:397)→ 实例化组件类(base.py:402-411)→ _build_results 真正执行组件(base.py:413)→ 校验并置 built = True(base.py:420-422)。组件实例化与 Component 类的关系见第 01 章。
3.6 第四步:结果怎么往下传 —— 是"拉",不是"推"
3.6.1 核心认知
调度器从不把上游的结果搬给下游。它只负责一件事:决定什么时候允许下游开始构建。真正的数据传递发生在下游构建的第一步——下游主动 await 上游要结果。
ChatInput Prompt
(已 built) (刚被允许构建)
│ │
│ │ ① _build_each_vertex_in_params_dict
│ │ 扫 raw_params,发现某个值是 Vertex 对象
│ ② await 上游.get_result(self, target_handle_name=key)
│◄───────────────────────────┤
│ │
│ ③ 按 handle 名取出对应 output 的值
├───────────────────────────►│
│ │ ④ params[key] = 真实值,然后才执行组件
为什么参数里会躺着 Vertex 对象?因为建图时,一条边的目标参数就是被直接填成源顶点对象的(Vertex._set_params_from_normal_edge,graph/vertex/base.py:308-334)。详见第 02 章。
3.6.2 替换的四种情形
_build_each_vertex_in_params_dict(graph/vertex/base.py:554)遍历 raw_params,按值的形状分四路:
| 值的形状 | 处理 | 位置 |
|---|---|---|
单个 Vertex | _build_vertex_and_update_params 拉一次结果 | graph/vertex/base.py:662 |
Vertex 列表 | _build_list_of_vertices_and_update_params 逐个拉,结果 合并成一个列表 | graph/vertex/base.py:670 |
dict 里含 Vertex | _build_dict_and_update_params 逐键拉 | graph/vertex/base.py:579 |
| 普通值 | 直接抄进 params | graph/vertex/base.py:572-573 |
两个细节值得单独记:
自引用会被丢弃。 if value == self: del self.params[key]; continue(graph/vertex/base.py:557-559)——顶点若在参数里引用了自己,直接删掉这个参数。考虑到 get_result 也要拿同一把顶点锁,这行同时避免了自死锁。
列表合并会跳过被条件剪掉的分支。 _build_list_of_vertices_and_update_params 里显式跳过「没 built 且在 conditionally_excluded_vertices 里」的顶点(graph/vertex/base.py:678-684),否则会往列表里塞进一个模板默认值(比如空 Message)当作真数据。源码注释把这个坑写得很清楚。
3.6.3 按 handle 名挑 output
get_result 是加锁的薄壳(graph/vertex/base.py:601-610),真逻辑在子类。
基类 Vertex._get_result(graph/vertex/base.py:644)很简单:没 built 就报错,built 了就按 use_result 返回 built_result 或 built_object。
ComponentVertex._get_result(graph/vertex/vertex_types.py:93)才是画布组件实际走的那条,它要解决一个真问题:一个组件可以有多个 output,下游到底要哪一个? 答案藏在边上—— 边记着源 handle 名和目标 handle 名。
# 真实源码节选 graph/vertex/vertex_types.py:129-137 —— 双向匹配 handle
edges = self.get_edge_with_target(requester.id)
for edge in edges:
if (
edge is not None
and edge.source_handle.name in self.results
and edge.target_handle.field_name == target_handle_name
):
匹配条件是两头都要对上:源 handle 名要在自己的 results 里,目标 handle 名要等于下游传来的 target_handle_name。命中后优先读 Output.value,只有当它是 UNDEFINED 时才回退到 results[source_handle.name](vertex_types.py:138-146)。取不到就一路抛带上下文的 ValueError(vertex_types.py:148-156)。
未构建时还有两条特殊出路(vertex_types.py:101-121):被条件剪掉的前驱返回下游那个输入的模板默认值;环边则按 edge.target_param 解析默认值。两者都属于第 04 章。
3.6.4 打包:finalize_build
组件跑完后,Vertex.build 调 finalize_build()(graph/vertex/base.py:872)把散落的东西装进一个 ResultData:
| 字段 | 内容 |
|---|---|
results | get_built_result() 归一化后的结果 |
artifacts | 展示用产物(文本/表格等) |
outputs | outputs_logs,每个 output 的日志 |
messages | 从 artifacts 里抽出的消息 |
token_usage | token 用量 |
component_display_name / component_id | 身份信息 |
基类版本在 graph/vertex/base.py:534,ComponentVertex 重写了一版(graph/vertex/vertex_types.py:194),差别在于消息从 result_dict 而非 artifacts_raw 里抽。这个 ResultData 就是前端拿到的那坨节点结果。
另有一条给 requester 的返回值路径:get_requester_result(graph/vertex/base.py:889)。没有 requester(调度器直接调 build 就是这种)返回 built_object;有 requester 则走边的 get_result_from_source(graph/edge/base.py:320),先 honor 履约再返回边上缓存的结果。
3.7 可回放的调试设施:快照与调用顺序
调度器出问题时最难受的是「跑完了才发现顺序不对,但现场没了」。Langflow 留了一套很轻的记录设施。
get_snapshot(graph/graph/base.py:2013)深拷贝六样东西:
| 快照字段 | 记的是 |
|---|---|
run_manager | 四个集合的当前状态(to_dict,runnable_vertices_manager.py:13) |
run_queue | 待跑队列 |
vertices_layers | 排序产出的剩余层 |
first_layer | 起跑名单 |
inactive_vertices / activated_vertices | 失活与激活标记 |
_record_snapshot(graph/graph/base.py:2025)把快照追加进 self._snapshots,并把 vertex id 追加进 self._call_order(两者初始化在 graph/graph/base.py:188-189)。触发点只有两处:prepare 结束时打一张空底(base.py:2308),以及每次 astep 结束时打一张带 vertex id 的(base.py:1640)。
于是 _call_order 就是这次运行的实际执行顺序,_snapshots[i] 是第 i 步之后的完整调度状态。测试和排障时,把两者一对,就能回答「为什么这个组件比那个先跑」。
注意两点局限:批处理路径(process)不打快照,所以这套设施只服务 astep 单步流;_snapshots 无上限、每张都是深拷贝,长流程会持续吃内存。
3.8 巧妙之处(可带走的技术)
- 排序只负责起跑,推进交给对账。 分层结果里只有
first_layer和vertices_to_run真正进了执行决策(§3.3.6)。这让调度天然容得下「跑到一半图变了」——条件剪枝、环回跳都只改集合,不必重排。 vertices_being_run的预占语义。 在锁里算出下一批就立刻塞进这个集合(graph/graph/base.py:2289-2293),后续任何并发的就绪判定都会在第二道闸被挡掉(runnable_vertices_manager.py:60-61)。用一个集合替掉了一整套分布式协调。- 前驱表反转出后继表。 只维护
run_predecessors一份真相,run_map由它一次性反转生成(runnable_vertices_manager.py:109-116)。销账只需查一次表。 - 卡住时回头捞前驱。
find_runnable_predecessors_for_successor(graph/graph/base.py:2837)递归回溯,把「后继没就绪」变成「那就找找它欠的谁能跑」,避免非首层顶点永远没人调度(§3.4.4)。 - 拉模型让缓存和剪枝都变简单。 因为值是下游主动拉的,frozen 复用(§3.5.4)、条件剪枝返回默认值(
vertex_types.py:101-109)都只需在「被拉的那一刻」做判断,不需要在推送侧铺一套旁路。 - 递归排序里的重入检测。
_sort_single_layer_by_dependency用processing集合识别层内环并就地切断(graph/graph/utils.py:781-786),把「可能无限递归」变成「返回当前下标」,不抛异常也不卡死。
3.9 边界与局限(诚实说)
| 局限 | 具体表现 | 依据 |
|---|---|---|
| 同层失败不会掐断同伴 | gather(return_exceptions=True) 已等全部结束,之后的 cancel() 基本是空操作 | graph/graph/base.py:2355、2000 |
| 并发粒度只有「一层」 | process 必须等整层最慢的那个跑完才开下一批,层内一个慢 LLM 会拖住所有人 | graph/graph/base.py:2239-2247 |
| 没有优先级/权重调度 | 有 sort_by_avg_build_time 的想法但从未接线 | graph/graph/base.py:2795(无调用者) |
| 首层顺序被字典序覆盖 | prepare 里 self._first_layer = sorted(first_layer),层内依赖排序的成果在首层被字典序盖掉;process 路径不经过 prepare,保留原序 | graph/graph/base.py:2683 vs 1826 |
| 快照只覆盖单步路径 | process 全程不调 _record_snapshot | graph/graph/base.py:2187-2249 |
| 死代码若干 | refine_layers、sort_interface_components_first、sort_by_avg_build_time、next_vertex_to_build 全仓库无调用者 | 见 §3.3.6 |
| 环的迭代上限是硬编码 | MAX_CYCLE_APPEARANCES = 2,模块级常量,不可配置 | graph/graph/utils.py:11 |
本章不覆盖: 环的判定与条件分支剪枝(cycle_vertices、conditionally_excluded_vertices、mark_branch)见第 04 章;HTTP 端点、SSE 事件与前端逐节点渲染见第 05 章。
3.10 代码地图(导航索引)
路径相对克隆根。
| 主题 | 文件 | 符号 |
|---|---|---|
| 排序总入口 | src/lfx/src/lfx/graph/graph/utils.py | get_sorted_vertices |
| 分层拓扑排序 | src/lfx/src/lfx/graph/graph/utils.py | layered_topological_sort、MAX_CYCLE_APPEARANCES |
| 层内排序 | src/lfx/src/lfx/graph/graph/utils.py | sort_chat_inputs_first、sort_layer_by_dependency、_sort_single_layer_by_dependency |
| 剪枝 | src/lfx/src/lfx/graph/graph/utils.py | filter_vertices_up_to_vertex、filter_vertices_from_vertex |
| 起点识别 / 迭代闸 | src/lfx/src/lfx/graph/graph/utils.py | find_start_component_id、should_continue |
| 未被调用的排序 helper | src/lfx/src/lfx/graph/graph/utils.py | refine_layers |
| 排序落地 + 准备 | src/lfx/src/lfx/graph/graph/base.py | Graph.sort_vertices、Graph.prepare、Graph.first_layer |
| 就绪判定(图侧) | src/lfx/src/lfx/graph/graph/base.py | Graph.is_vertex_runnable、Graph.build_run_map |
| 找下一批 | src/lfx/src/lfx/graph/graph/base.py | find_next_runnable_vertices、find_runnable_predecessors_for_successor、get_next_runnable_vertices |
| 就绪判定(管理器) | src/lfx/src/lfx/graph/graph/runnable_vertices_manager.py | RunnableVerticesManager、is_vertex_runnable、are_all_predecessors_fulfilled、build_run_map、remove_from_predecessors |
| 单步执行 | src/lfx/src/lfx/graph/graph/base.py | Graph.astep、get_next_in_queue、extend_run_queue、async_start |
| 批并发执行 | src/lfx/src/lfx/graph/graph/base.py | Graph.process、Graph._execute_tasks、Graph._run、Graph.arun |
| 顶点构建 + frozen 缓存 | src/lfx/src/lfx/graph/graph/base.py | Graph.build_vertex、VertexBuildResult |
| 顶点内部构建 | src/lfx/src/lfx/graph/vertex/base.py | Vertex.build、Vertex._build、Vertex._build_results、Vertex.steps |
| 参数里的上游替换 | src/lfx/src/lfx/graph/vertex/base.py | _build_each_vertex_in_params_dict、_build_vertex_and_update_params、_build_list_of_vertices_and_update_params |
| 结果拉取 | src/lfx/src/lfx/graph/vertex/base.py | Vertex.get_result、Vertex._get_result、get_requester_result |
| 按 handle 挑 output | src/lfx/src/lfx/graph/vertex/vertex_types.py | ComponentVertex._get_result、get_edge_with_target |
| 结果打包 | src/lfx/src/lfx/graph/vertex/base.py / vertex_types.py | finalize_build、ResultData |
| 边上的结果履约 | src/lfx/src/lfx/graph/edge/base.py | get_result_from_source |
| 调试快照 | src/lfx/src/lfx/graph/graph/base.py | Graph.get_snapshot、Graph._record_snapshot、_call_order |