跳到主要内容

数据截至 (上游 commit 7a975c596eca)

运行时与事件流 — 点一下 Run 之后前后端之间发生了什么

30 秒导读: 在 Langflow 画布上点 Run,浏览器并不是发一个请求等一个结果。它先发一个 POST 拿到 job_id,再开一条长连接把后端吐出来的事件流一条条读回来;后端一边跑图一边把「排序完成 / 某个点开始 / 某个点结束 / 又来一个 token / 结束」写进一个队列,由这条连接推给前端。本章把这条链路从按钮一路追到 asyncio.Queue,再追回节点变色。

本章不讲组件怎么被发现、怎么变成 MCP 工具(那是 06),也不细讲"谁先跑、谁能并行"的调度规则(那是 03)。这里只关心一次交互式运行的端到端链路和事件协议


1. 先看全貌:一次 Run 的三段旅程

一次交互式运行被切成三个独立的 HTTP 交互,而不是一个。

浏览器 后端
│ │
│ ① POST /api/v1/build/{flow_id}/flow│
├───────────────────────────────────>│ 建 job_id + asyncio.Queue
│ │ 起后台 task 跑 generate_flow_events
│ <──────── {"job_id": "..."} ───────┤ │
│ │ │ 边跑边把事件塞进队列
│ ② GET /api/v1/build/{job_id}/events│ ▼
├───────────────────────────────────>│ ┌───────────────┐
│ <═══ NDJSON 事件流(长连接) ═════════┤ │ asyncio.Queue │
│ vertices_sorted / build_start │ └───────────────┘
│ end_vertex / token / end │
│ │
│ ③ POST /api/v1/build/{job_id}/cancel(用户点停止时)
├───────────────────────────────────>│ event_task.cancel()

三段各自的职责:

阶段路由干什么定义在哪
① 建 jobPOST /build/{flow_id}/flow建队列、起后台任务、立刻返回 job_id,不等结果src/backend/base/langflow/api/v1/chat.py:256 (build_flow)
② 拉事件GET /build/{job_id}/events把队列里的事件按投递方式送给浏览器src/backend/base/langflow/api/v1/chat.py:414 (get_build_events)
③ 取消POST /build/{job_id}/cancel取消那个后台任务src/backend/base/langflow/api/v1/chat.py:441 (cancel_build)

三个路由的实现全部落在同一个文件里:src/backend/base/langflow/api/build.pystart_flow_buildget_flow_events_responsecancel_flow_build

为什么要拆成两个请求? 因为「启动」和「读结果」的失败语义完全不同。POST 失败意味着这张流根本没跑起来(校验不过、组件被禁用),要立刻报错;GET 失败只意味着这条连接断了——任务还在后端跑,重连就能接着读。拆开之后,重连不需要重跑。


2. 第一段:POST 只做一件事——建 job

start_flow_build 短得出奇,因为它真正只做三件事。

job_id = str(uuid.uuid4())
_, event_manager = queue_service.create_queue(job_id) # 1. 建队列 + 事件管理器
task_coro = generate_flow_events(...) # 2. 准备协程(还没跑)
queue_service.start_job(job_id, task_coro) # 3. 丢进后台
return job_id

src/backend/base/langflow/api/build.py:182-250,符号 start_flow_build

三件事对应到 JobQueueService 的两个方法:

  • create_queuesrc/backend/base/langflow/services/job_queue/service.py:178)建一个无界 asyncio.Queue,并用它构造一个 EventManager,登记进 self._queues[job_id]
  • start_job(同文件 :206)用 asyncio.create_task 把协程跑起来,任务句柄存回同一条记录。

start_job 有个不显眼但关键的包装:协程被 _guarded_task 裹了一层(:246)。它的作用是——万一 generate_flow_events 抛了没接住的异常,也要保证队列里先落一条 error 事件、再落一个结束哨兵 (None, None, ts) 才退出。没有这层包装,消费端会永远卡在 queue.get() 上等一个永远不来的结束信号。

顺带澄清一个容易混淆的文件名:src/backend/base/langflow/api/v1/flow_events.py 不是本章说的作业队列。那是另一套东西——GET/POST /api/v1/flows/{flow_id}/events,给某张流追加/读取"设计期通知事件"(FlowEventsServicesrc/backend/base/langflow/services/flow_events/service.py:44),和运行时的构建事件流没有关系。运行时的作业队列在 services/job_queue/service.py


3. 第二段:同一批事件,三种投递方式

get_flow_events_responsebuild.py:253)根据查询参数 event_delivery 分岔。三种模式的取舍:

模式传输形态谁在用代价
流式streaming一次 GET,长连接,NDJSON 边算边推前端默认需要中间层不缓冲响应
轮询polling反复 GET,每次把队列排空返回,25ms 后再来流式不可用时的降级延迟 = 轮询间隔
直连directPOST /build 本身就返回流,不给 job_id前端 e2e 测试无法重连

枚举定义在 src/backend/base/langflow/api/utils/core.py:85EventDeliveryType),前端同名枚举在 src/frontend/src/constants/enums.ts:42

┌── streaming ─ 一次 GET,长连接,事件产生即推
生产端 asyncio.Queue ────┼── polling ─── 反复 GET,每次 drain 光队列
└── direct ──── POST 直接变成流

3.1 流式:create_flow_response 的三个零件

create_flow_responsebuild.py:342)返回一个 DisconnectHandlerStreamingResponse,内部有三个协作的零件。

零件一,消费循环 consume_and_yieldbuild.py:392-407)。就是 await queue.get() 然后 yield _project_event_to_v1(value.decode("utf-8")),直到取到 value is None 这个哨兵才 break。事件在队列里放了多久也顺手记进 debug 日志。

零件二,心跳 _heartbeatbuild.py:367-382)。每 STREAMING_ACTIVITY_REFRESH_S(默认 10 秒,build.py:60)调一次 queue_service.touch_activity(job_id)。注释里点明了它为什么必须独立于事件节奏:一次慢 LLM 调用可能几十秒不产事件,如果用"最近一次事件时间"来判活,看门狗会把一个正常跑着的 job 当成僵尸回收。心跳把"客户端还在"和"任务还在产出"这两件事解耦了。

心跳只在后端支持时才起:touch = getattr(queue_service, "touch_activity", None),内存版 JobQueueService 根本没这个方法,于是 heartbeat_taskNone,整段逻辑自动消失(build.py:365,384)。touch_activity 只有 Redis 版实现(services/job_queue/service.py:1081)。

零件三,断线处理 on_disconnectbuild.py:409-430)。浏览器一关,DisconnectHandlerStreamingResponse.listen_for_disconnectsrc/backend/base/langflow/api/disconnect.py)收到 http.disconnect,回调就会:停心跳 → event_task.cancel() 取消构建 → 清队列 → 发一个 on_end用户关标签页,后端就真的停算了,不会把一次没人要的 LLM 调用跑完。

多进程部署下还有一条支线:如果这条 GET 落在了没拥有那个 task 的 worker 上(event_task is None),就退而求其次调 queue_service.signal_cancel(job_id),通过 Redis 频道把取消信号发给真正的 owner(build.py:414-424)。

3.2 轮询:一次请求排空一次队列

轮询分支(build.py:295-330)的做法很直白:

while not main_queue.empty(): # 有多少拿多少,不阻塞
_, value, _ = await main_queue.get()
if value is None: ... # 哨兵 → 取消 event_task、发 on_end
events.append(value.decode("utf-8"))
if not events: # 一条都没有 → 阻塞等第一条,避免空转
_, value, _ = await main_queue.get()
return Response(content="\n".join(...), media_type="application/x-ndjson")

注意 if not events 那一支:如果这一轮什么都没有,它会阻塞等一条再返回。这让"没有事件"的空转成本被摊薄到服务端的一次 await,而不是客户端疯狂空跑。

event_delivery 落到未知值时不会静默按轮询处理,而是直接 400 并把支持的值列出来(build.py:277-293)——注释说明了理由:三种模式的跨 worker 保障不同,静默 fallthrough 会掩盖多 worker 配置错误。


4. 第三段:后端到底在跑什么

generate_flow_eventsbuild.py:439)是整台机器的主体。它内部定义了一串闭包,链路是:

generate_flow_events

├─ build_graph_and_get_order() 建图 + 排序
│ ├─ create_graph() build_graph_from_db / build_graph_from_data
│ └─ sort_vertices() graph.sort_vertices(stop_id, start_id)

├─ on_vertices_sorted ← 第一条事件:告诉前端「这些点要跑」
├─ on_build_start

└─ _run_vertex_build()
└─ 对第一层每个 id: create_task(build_vertices(id))

└─ _build_vertex → graph.build_vertex
然后 on_end_vertex
然后对 next_vertices_ids 递归 create_task

4.1 建图与排序

create_graphbuild.py:615-659)二选一:

  • 请求体没带 data → 从数据库读那张流:build_graph_from_dbbuild.py:625)。
  • 带了 data(画布上未保存的改动)→ 直接用请求体建:build_graph_from_databuild.py:644)。

这就是"画布上改了没保存也能点 Run"的实现。图怎么从 JSON 变出来,见 02-graph-construction

sort_verticesbuild.py:661-666)调 graph.sort_vertices(stop_component_id, start_component_id),拿到第一层可跑的点。它外面套了一层 try/except:带 start/stop 的排序失败时,退回无参数的全图排序,而不是让整次运行崩掉。

排完序立刻发第一条事件(build.py:928):

event_manager.on_vertices_sorted(data={"ids": ids, "to_run": vertices_to_run})

ids 是第一层,to_run 是整次运行会碰到的全部点——前端用后者把所有相关节点先标成 TO_BUILD(灰色待跑),用前者标成"马上开始"。

4.2 并发是"铺开"出来的,不是排好的

这是本章最值得带走的一处设计。build_verticesbuild.py:833-883)结构如下:

vertex_build_response = await _build_vertex(vertex_id, graph, event_manager) # 跑这个点
...
event_manager.on_end_vertex(data={"build_data": build_data}) # 报告结果

if vertex_build_response.valid and vertex_build_response.next_vertices_ids:
tasks = [asyncio.create_task(build_vertices(nid, ...)) # 对每个后继开一个 task
for nid in vertex_build_response.next_vertices_ids]
await asyncio.gather(*tasks)

next_vertices_ids 来自 graph.get_next_runnable_vertices(lock, vertex=vertex, cache=False)build.py:697)。所以并发度不是谁事先算好的,而是每个点跑完自己报告"接下来谁能跑",然后就地为每个后继开一个 task,并发自然铺开:

第一层 ids = [A, B]
├─ task(A) ── 完 ── next=[C] ──── task(C) ── 完 ── next=[E] ── task(E)
└─ task(B) ── 完 ── next=[C,D] ─┬ task(C) (已在跑的会被 run_manager 挡掉)
└ task(D)

"谁算可跑、重复的怎么去重、环怎么办"全部封在 graph.get_next_runnable_verticesRunManager 里——见 03-scheduler04-cycles-and-branching。运行时这一层只负责照单开 task

对比一下同一个引擎的另一种驱动方式会更清楚。非交互路径用的 Graph.processsrc/lfx/src/lfx/graph/graph/base.py:2187-2248)是层同步的:一层的 task 全 gather 完,_execute_tasks 返回下一层,才开下一批。两种驱动的差别:

/build(交互)Graph.process/run
驱动方式每点跑完就地递归 spawn攒成一层,gather 完再开下一层
层间屏障
后果快的分支不等慢的分支一层里最慢的点拖住整层
代码build.py:867-883graph/graph/base.py:2216-2248

4.3 单点执行 _build_vertex

_build_vertexbuild.py:672-831)是"跑一个点"的完整包裹,核心一句是 await graph.build_vertex(...)build.py:684),把 event_manager 一路传进去(graph/graph/base.py:2049)。它额外做的事:

  • 错误不炸流程:组件抛异常时不往上抛,而是把 traceback 包成 OutputValue(type="error") 塞进 result_data_responsevalid=Falsebuild.py:705-732)。因为这个点失败了,后继不该跑,但事件流必须继续,前端才能把这个节点涂红。
  • 落库与缓存:不流式且 log_builds 为真(且 graph.persist_messages 为真——匿名运行会把它置假、整体跳过落库)→ 后台任务 log_vertex_build 落库;否则把整张图写进 chat 缓存(build.py:743-757)。
  • 停点裁剪:如果用户指定了停在某个组件,且它就在后继里,就把后继裁成只剩它(build.py:777-778)。
  • 计时duration / timedelta 写回响应,前端拿它显示每个节点跑了多久。

最后收尾(build.py:1007-1009):把所有点的耗时求和,发 on_end(data={"build_duration": ...}),再往队列里放结束哨兵 (None, None, time.time())build.py:1032)。哨兵是流结束的唯一信号,前面提到的 _guarded_task 就是为了保证异常路径上它也一定会被放进去。


5. 事件对象:EventManager 与它的三层容错

前面反复出现的 event_manager.on_xxx(...) 到底是什么?EventManagersrc/lfx/src/lfx/events/event_manager.py:30)是一个把"发生了什么"翻译成队列里一行 JSON 的适配器——只有大约 120 行,但每一处容错都有故事。

5.1 事件清单

create_default_event_managerevent_manager.py:124-136)登记了全部 10 种事件。这就是前后端之间的完整协议:

方法名线上事件名谁发前端拿它干什么
on_vertices_sortedvertices_sortedbuild.py:928把待跑节点标成 TO_BUILD
on_build_startbuild_startbuild.py:931 / 组件id 时记录起跑时刻;有 id 时标 BUILDING
on_end_vertexend_vertexbuild.py:867主力事件:写结果、涂色、点亮下一批边
on_build_endbuild_end组件BUILT
on_tokentoken组件流式输出往聊天气泡追加字符
on_messageadd_message组件新增一条聊天消息
on_remove_messageremove_message组件撤回一条消息
on_loglogcustom_component/component.py:1782追加日志到 flowPool
on_errorerrorbuild.py:899,969弹错误、标红
on_endendbuild.py:1008收尾:算总时长、isBuilding=false

前端 onEventswitch 分支(src/frontend/src/utils/buildUtils.ts:721-849)和这张表一一对应,没有多余分支也没有遗漏。

5.2 注册:register_event 用 partial 把 event_type 焊死

def register_event(self, name, event_type, callback=None):
if not name.startswith("on_"): raise ValueError(...)
callback_ = partial(self.send_event, event_type=event_type) # 把类型名预先绑上
self.events[name] = callback_

event_manager.py:54-70

调用方于是只需要写 on_end_vertex(data=...),不必每次重复事件类型字符串。名字必须以 on_ 开头是硬约束。

5.3 发送:send_event 的三层容错

send_eventevent_manager.py:72-115)短短 40 行里叠了三层"绝不因为发事件而搞崩流程"的保护。

第一层:序列化失败就降级,不抛。

try:
jsonable_data = jsonable_encoder(data)
except (ValueError, TypeError) as exc:
jsonable_data = serialize(data, to_str=True) # 兜底:不认识的对象转字符串

event_manager.py:73-84。注释点名了触发场景(issue #12591):某个向量库客户端里揣着一个 threading.Lock,通过 Loop 组件的每轮构建事件冒出来,FastAPI 的 jsonable_encoder 直接抛 ValueError如果不降级,整次流程会因为"一个日志字段没法转 JSON"而挂掉。

第二层:跨线程用 call_soon_threadsafe

if in_event_loop:
self.queue.put_nowait(item)
elif self._loop is not None and self._loop.is_running():
self._loop.call_soon_threadsafe(self.queue.put_nowait, item)
else:
self.queue.put_nowait(item) # 无 loop 的同步环境,如单测

event_manager.py:97-111self._loop 在构造时抓一次(event_manager.py:35)。

这不是防御性编程的空转——真有调用方在别的线程里。token 事件就是:Component._process_chunkawait asyncio.to_thread(self._event_manager.on_token, ...) 把它扔到线程池里发(src/lfx/src/lfx/custom/custom_component/component.py:2122-2128)。直接 put_nowait 虽然也能进队列,但唤不醒在 queue.get() 上等着的那个协程。

第三层:队列满只丢事件,不炸流程。

except asyncio.QueueFull:
logger.warning("Event queue full; dropping event_type=%s", event_type)

event_manager.py:112-113。取舍很明确:丢一条 token 的可观测性 ≪ 让整次运行失败

5.4 兜底:没注册的事件是个 noop

def noop(self, *, data): pass
def __getattr__(self, name): return self.events.get(name, self.noop)

event_manager.py:117-121

这行代码解释了 lfx 为什么能脱离 Web 界面独立运行。create_stream_tokens_event_managerevent_manager.py:139)比默认版少注册了 on_vertices_sortedon_remove_message;组件代码里照样敢直接写 event_manager.on_vertices_sorted(...)——落到 __getattr__ 变成 noop,什么也不发生。组件不需要知道自己跑在哪种运行时里。

下面这段示意代码把这个模式提炼出来(# 示意,非源码):

# 示意,非源码:为什么 noop 兜底能让同一份组件代码跑在不同运行时里
class Emitter:
def __init__(self, registered):
self.events = registered # 这个运行时支持哪些事件

def _noop(self, *, data): pass

def __getattr__(self, name): # 没登记的事件名 → 静默丢弃
return self.events.get(name, self._noop)

em = Emitter({"on_token": print})
em.on_token(data="hi") # 有登记 → 真的发
em.on_vertices_sorted(data=[]) # 没登记 → 什么都不做,也不报错

重点看:调用方永远不用先判断"这个运行时支不支持这个事件"。


6. 构建结果的三个去向

同一次点构建,结果会分头去三个地方。别把它们搞混:

去向函数目的位置
给浏览器event_manager.on_end_vertex实时更新画布src/backend/base/langflow/api/build.py:867
给数据库log_vertex_build事后查历史构建src/lfx/src/lfx/graph/utils.py:348
给 webhook 订阅者emit_vertex_build_eventwebhook 触发时也让 UI 有动画src/lfx/src/lfx/graph/utils.py:122

log_vertex_build 是纯落库,且有两级降级:设置里没开 vertex_builds_storage_enabled 直接返回;开了 telemetry writer 就把行丢进它的磁盘 outbox(graph/utils.py:413-416),避开请求线程的数据库连接池;writer 没跑起来才直接写库。最外层还包了一个 except Exception: logger.warninggraph/utils.py:441)——记录构建历史失败,绝不影响构建本身

emit_vertex_build_event 走的是另一条完全独立的通道:webhook_event_manager。第一件事就是 if not webhook_event_manager.has_listeners(flow_id_str): returngraph/utils.py:156)——没人订阅就一个字节都不生产。它由 Graph._execute_tasks 在算出 next_runnable_vertices 之后才调用(graph/graph/base.py:2412),因为 payload 里需要 next_vertices_idsgraph/utils.py:425-427 的注释专门解释了为什么不能在 log_vertex_build 里顺手发。

它也有个不显眼的 except ImportError: passgraph/utils.py:202):lfx 可以在没装 langflow 的环境里独立使用,那时 webhook_event_manager 根本不存在。


7. 前端:从字节流到节点变色

链路的另一半在 src/frontend/src/utils/buildUtils.ts

7.1 主入口与降级

buildFlowVerticesWithFallbackbuildUtils.ts:160-181)只做一件事:先按配置的投递方式试,撞上"流式不可用"就换轮询重来一次。

try {
return await buildFlowVertices({ ...params });
} catch (e) {
if (e.message === POLLING_MESSAGES.ENDPOINT_NOT_AVAILABLE ||
e.message === POLLING_MESSAGES.STREAMING_NOT_SUPPORTED) {
return await buildFlowVertices({ ...params, eventDelivery: EventDeliveryType.POLLING });
}
throw e;
}

两个哨兵消息定义在 src/frontend/src/constants/constants.ts:935-938。这层降级是给那些会缓冲响应体、把长连接一次性吞掉的反向代理准备的。

buildFlowVerticesbuildUtils.ts:210)本体按模式分三条路:

  • direct:直接对 POST /build 做流式读(buildUtils.ts:282-328),不拿 job_id
  • streaming:先 fetch(buildUrl)job_idbuildUtils.ts:336-361),再对事件 URL 做流式读(:385-421)。
  • polling:拿到 job_id 后交给 pollBuildEvents:431)。

轮询实现在 src/frontend/src/customization/utils/custom-poll-build-events.ts:循环 GET,把 NDJSON 按行 JSON.parse单行解析失败只跳过这一行:66-68),看到 end 事件就停,否则 BUILD_POLLING_INTERVAL(25ms,constants.ts:955)后再来。

7.2 解帧:按 \n\n 切,跨 chunk 拼

performStreamingRequestsrc/frontend/src/controllers/API/api.tsx:330)负责把字节流切成事件。后端每条事件是 json.dumps(...) + "\n\n"event_manager.py:87),前端就按 \n\n 切(api.tsx:377)。

TCP chunk 边界和事件边界不对齐,所以有个 current: string[] 缓冲:一段字符串如果不是以 } 结尾、或者拼起来 JSON.parse 还是失败,就先攒着,等下一个 chunk 续上(api.tsx:381-394)。

请求头里那句 Connection: "close" 带着注释——"this flag is fundamental to ensure server stops tasks when client disconnects"(api.tsx:343)。它是第 3.1 节 on_disconnect 那套取消机制在前端这一侧的前提。

7.3 批处理:为什么高频 token 不会卡死 React

这是前端最值得看的一处工程决策。

问题:一次流式对话每秒能产上百条 token 事件,每条都同步 set() 一次 Zustand,React 就要重渲染上百次,界面直接卡住。

做法:把事件分成两类。

export const BATCHABLE_EVENTS = new Set(["end_vertex", "build_start", "build_end"]);
export const BATCH_YIELD_MS = 50;

buildUtils.ts:59-65

processBatchedEventsbuildUtils.ts:606-668)的处理逻辑:

  • 属于 BATCHABLE_EVENTS 的 → 同步处理,不 await。因为一旦 await,就跨过了微任务边界,React 18 的自动批处理就断了,同一个 chunk 里的多次 set() 会变成多次渲染。
  • 不属于的(token / add_message / error / log…)→ 老老实实 await onEventFallback(...),它们有真正的异步副作用。
  • 这批里只要有过 batchable 事件,最后 await new Promise(r => setTimeout(r, BATCH_YIELD_MS))buildUtils.ts:665主动让出 50ms,给浏览器一次真正提交渲染、响应交互的机会。

调度侧的配合在 api.tsx:398-403:一个 chunk 解出来的全部事件先攒成数组,整批交给 onDataBatch,而不是一条条 onData

// 示意,非源码:批处理与逐条处理的渲染次数差别
// 逐条:await 每条 → 每条一次 set → N 次渲染
for (const e of events) { await handle(e); }

// 成批:同类事件同步处理 → React 合并成 1 次渲染 → 再让出 50ms
for (const e of events) { if (batchable(e)) handleSync(e); else await handle(e); }
await sleep(50);

重点看:同步与否决定了 React 会不会把多次状态更新合并成一次渲染。

7.4 end_vertex:一条事件干了四件事

processEndVertexEventbuildUtils.ts:456-596)被单独抽成同步函数,注释写明了原因:轮询模式要能不 await 地调它,否则批处理失效。它依次做:

  1. 判定成败buildData.valid 为假就从 outputs 里把错误信息掏出来,调 onBuildError,状态置 ERROR;否则置 BUILT:474-503)。
  2. 算分段耗时:如果这个点是输出类节点,用 Date.now() - buildStartTime 算出这一段的耗时,写进最后一条机器消息的 build_duration,同时更新 React Query 缓存、Zustand、以及持久化到后端(:505-567)。
  3. 点亮下一批边clearAndSetEdgesRunning(nextIds) 一次遍历里把所有边的动画清掉,再把源头在 nextIds 里的边设成 animated + className:"running":575,实现在 src/frontend/src/stores/flowStore.ts:1196-1209)。这就是画布上那圈流动的虚线。
  4. 推进层次:把 next_vertices_ids 标成 TO_BUILD 并回调 onBuildStart:581-593)。

7.5 状态怎么落到节点上

useFlowStoresrc/frontend/src/stores/flowStore.ts)持有构建状态。关键几处:

做什么符号位置
发起构建、装配所有回调buildFlowflowStore.ts:810
end_vertex 的主回调handleBuildUpdateflowStore.ts:963
节点 id → 构建状态 + 时间戳updateBuildStatusflowStore.ts:1246
节点 id → 构建产物(结果/日志)addDataToFlowPoolflowStore.ts:234
边的动画开关clearAndSetEdgesRunningflowStore.ts:1196
收尾把残留的 BUILDING 改回 BUILTrevertBuiltStatusFromBuildingflowStore.ts:1259

handleBuildUpdate 里有个容易忽略的动作:把 next_vertices_idstop_level_verticeszip 再过滤掉已在上一层出现过的,才追加成新的一层(flowStore.ts:989-1015)。因为后端是"每个点各自报告后继",同一个后继会被多个前驱各报一次,前端必须自己去重,否则层列表会膨胀。

7.6 旧路径:updateVerticesOrder

updateVerticesOrderbuildUtils.ts:99-158)先调 getVerticesOrder 拿到排序,再自己一层层 POST 每个点——就是同文件里那个 buildVerticesbuildUtils.ts:853)配套的旧模型。它现在被事件流路径取代了:新路径里同样的信息由 vertices_sorted 事件送达(buildUtils.ts:722-748),少一次往返。对应的后端路由 POST /build/{flow_id}/vertices/{vertex_id} 已标 deprecated=True, include_in_schema=Falsechat.py:460)。


8. 取消与断线:三条路径通向同一个 cancel

触发前端动作后端动作
用户点停止buildController.abort() → abort 监听器 POST /build/{job}/cancelcancel_flow_buildqueue_service.cancel_job
关标签页 / 断网连接断开on_disconnectevent_task.cancel()
某个组件构建失败onBuildErrorbuildController.abort()同第一行

前端把取消绑在 AbortController 的 abort 事件上(buildUtils.ts:366-379):一次 abort 既停掉本地的读流,又顺手发出取消请求。

cancel_flow_buildbuild.py:1035-1117)的返回值语义值得留意——它区分了三种"没取消成":

  • event_task.done() → 返回 True没得取消也算成功build.py:1089-1091)。
  • event_task is None 且后端支持跨 worker 取消 → 通过 Redis 发信号,返回 Truebuild.py:1064-1078)。
  • event_task is None 且不支持 → 返回 False,注释诚实承认:"任务已经跑完被清理了"和"任务在一个够不着的 worker 上"这两种情况无法廉价区分build.py:1079-1087)。

取消路径上还有一处细节:_run_vertex_build 捕到 CancelledError 时,用 asyncio.create_task 而不是 background_tasks.add_task 去收尾 trace(build.py:945-952)。原因写在注释里——POST /build 的响应早就返回了,FastAPI 的 background_tasks 队列已经排干,这时 add_task 会被静默丢弃


9. 非交互入口:/run 和 webhook 有什么不同

同一台图执行引擎,还有两个不经过 job 队列的入口。

POST /build/{flow_id}/flowPOST /run/{flow_id_or_name}POST /webhook/{flow_id_or_name}
谁在用画布 / Playground外部程序、SDK外部系统推事件
认证会话用户API keywebhook 认证
返回job_id,结果走事件流RunResponse(全部输出)或 SSE立刻 202,什么都不返回
驱动方式递归 spawn(4.2 节)Graph.arunprocess 层批/run
事件10 种全套stream=true 时 4 种仅当 UI 开着 SSE 时
入口api/v1/chat.py:256api/v1/endpoints.py:962api/v1/endpoints.py:1194

/runsimplified_run_flowendpoints.py:962)走 _run_flow_internalendpoints.py:798)。分两支:

  • stream=false(默认)→ await simple_run_flow(...),一路 Graph.from_payloadgraph.set_run_idrun_graph_internalGraph.arunsrc/lfx/src/lfx/graph/graph/base.py:1212),同步等到全部跑完,返回 RunResponse(outputs=..., session_id=...)endpoints.py:494)。
  • stream=true → 自建一个 asyncio.Queuecreate_stream_tokens_event_manager,把 run_flow_generator 丢进 create_task,返回 text/event-streamendpoints.py:858-881)。注意它用的是精简版事件管理器(event_manager.py:139),只有 add_message / token / end / end_vertex / error / build_start / build_end / log,没有 vertices_sorted——因为没有画布要涂色。

webhookwebhook_run_flowendpoints.py:1194)更彻底:把请求体塞成 webhook 组件的 tweak,asyncio.create_task 起一个后台任务,立刻返回 202 {"status": "in progress"}endpoints.py:1262-1281)。调用方拿不到任何结果。

那 webhook 触发时,开着画布的用户为什么还能看到节点动?靠一条独立的 SSE 通道:GET /webhook-events/{flow_id_or_name}endpoints.py:1115)订阅 webhook_event_manager,超时就发 heartbeat(endpoints.py:1161-1175)。webhook 任务只在 has_ui_listeners 为真时才 emit_events=Trueendpoints.py:1257,1271)——没人看就不生产事件


10. 边界与坑

  • event_delivery 的默认值前后端不一致。 后端 build_flow 签名里默认 POLLINGchat.py:266),get_build_events 默认 STREAMINGchat.py:400),前端 buildFlow 默认 STREAMINGflowStore.ts:818)。实际生效的是前端显式带上的那个查询参数。
  • job_id 只是内存里的键。 内存版 JobQueueService 把队列存在进程字典里(services/job_queue/service.py:97)。多 worker 且没配 Redis 后端时,事件 GET 落到别的 worker 上就是 404。跨 worker 能力靠 RedisJobQueueServiceservice.py:719)。
  • 构建队列是无界的。 create_queueasyncio.Queue() 没有 maxsizeservice.py:198),所以 send_eventQueueFull 分支在默认配置下不会触发——它是为带背压的后端准备的。
  • 一次未被消费的构建也会跑完。 只要没人断线、没人取消,generate_flow_events 就一直跑,事件在队列里堆着。清理靠 JobQueueService 的周期清理 + 5 分钟宽限期(service.py:103)。
  • 轮询模式下 token 的实时感取决于轮询间隔。 25ms 已经很密,但和流式推送不是一回事。
  • 代码里看不出每种事件的正式 schema 定义——事件 payload 是各调用点手写的 dict,没有集中的 pydantic 模型约束。想知道某个事件带什么字段,只能去发送点看。

11. 代码地图

主题文件路径符号名
建 job、返回 job_idsrc/backend/base/langflow/api/build.pystart_flow_build
事件投递分岔(流式/轮询/直连)src/backend/base/langflow/api/build.pyget_flow_events_response
流式响应 + 心跳 + 断线取消src/backend/base/langflow/api/build.pycreate_flow_response, _heartbeat, on_disconnect
构建主流程src/backend/base/langflow/api/build.pygenerate_flow_events
建图与排序src/backend/base/langflow/api/build.pybuild_graph_and_get_order, create_graph, sort_vertices
单点执行与错误包裹src/backend/base/langflow/api/build.py_build_vertex
递归 spawn 并发src/backend/base/langflow/api/build.pybuild_vertices, _run_vertex_build
取消构建src/backend/base/langflow/api/build.pycancel_flow_build
断线检测的响应类src/backend/base/langflow/api/disconnect.pyDisconnectHandlerStreamingResponse
三个路由src/backend/base/langflow/api/v1/chat.pybuild_flow, get_build_events, cancel_build
队列服务与崩溃兜底src/backend/base/langflow/services/job_queue/service.pycreate_queue, start_job, _guarded_task
事件对象与三层容错src/lfx/src/lfx/events/event_manager.pyEventManager, send_event, register_event, __getattr__
事件清单src/lfx/src/lfx/events/event_manager.pycreate_default_event_manager, create_stream_tokens_event_manager
构建记录落库src/lfx/src/lfx/graph/utils.pylog_vertex_build
webhook 实时事件src/lfx/src/lfx/graph/utils.pyemit_vertex_build_event, emit_build_start_event
单点构建 / 层批驱动src/lfx/src/lfx/graph/graph/base.pybuild_vertex, process, _execute_tasks, arun
token 跨线程发送src/lfx/src/lfx/custom/custom_component/component.py_process_chunk, set_event_manager
前端主入口与降级src/frontend/src/utils/buildUtils.tsbuildFlowVerticesWithFallback, buildFlowVertices
事件分发src/frontend/src/utils/buildUtils.tsonEvent, processEndVertexEvent
批处理与让出src/frontend/src/utils/buildUtils.tsprocessBatchedEvents, BATCHABLE_EVENTS, BATCH_YIELD_MS
流解帧src/frontend/src/controllers/API/api.tsxperformStreamingRequest
轮询循环src/frontend/src/customization/utils/custom-poll-build-events.tscustomPollBuildEvents
构建状态落到画布src/frontend/src/stores/flowStore.tsbuildFlow, handleBuildUpdate, updateBuildStatus, clearAndSetEdgesRunning
非交互入口src/backend/base/langflow/api/v1/endpoints.pysimplified_run_flow, simple_run_flow, webhook_run_flow, webhook_events_stream

继续读: 图是怎么建出来的看 02-graph-constructionget_next_runnable_vertices 凭什么说某个点可跑,看 03-scheduler;后继里出现环怎么办,看 04-cycles-and-branching;一条流怎么变成别人能调的工具,看 06-ecosystem-and-exits