Workflow 总览:拖出来的一张流程图,是怎么被跑起来的
阶段2 讲清了"单次模型调用",但真实应用往往是一连串步骤:先检索知识库 → 再调 LLM → 再判断分支 → 再调工具。用户在画布上拖出这张"节点图",后端怎么执行?主角是 core/workflow/workflow_entry.py(入口)和 graph_topology.py(拓扑解析)。弄清三件事:①一个工作流由哪几样东西组成(图配置、Graph 对象、变量池、运行时状态);②WorkflowEntry 怎么把这些组装进 GraphEngine 并挂上各种"层"(限流/日志/观测);③执行为什么是"吐出一串事件"而不是"返回一个结果"。今天先建立整体世界观,Day09 钻节点、Day10 钻引擎。
GraphEngine 就是那列按线路行驶的地铁,WorkflowEntry 是发车前把线路图、时刻表、乘客都装配好的调度室。类比二:执行工作流像看直播而不是看录播——引擎不是跑完给你一个大结果,而是"每到一站就播报一次"(节点开始了、节点产出了、出错了)。run() 返回的是一个源源不断的事件流(generator),前端才能实时看到"正在检索…正在生成…"。痛点:画布上拖出来的图,后端怎么执行
{nodes:[...], edges:[...]})。后端拿到这份 JSON 后要解决一堆问题:从哪个节点开始跑?某节点的输入要用到前面节点的输出,数据怎么传递?多个节点能不能并行?跑到一半用户点了停止怎么办?执行进度怎么实时推给前端?跑太久/步数太多怎么熔断?如果这些逻辑和"具体某个节点干啥"混在一起,代码会彻底失控。graphon.graph_engine.GraphEngine)+ Dify 侧的装配入口 WorkflowEntry。引擎只懂"图论":按边推进、并行调度、事件广播、熔断控制;至于"某个节点具体干什么",引擎一概不管,交给节点自己(Day09)。WorkflowEntry 的活是"装配"——把图、变量池、运行时状态、各种控制"层"拼到引擎上,然后按一个 run() 键。一张图由哪几样东西组成
看 WorkflowEntry.__init__ 的入参(api/core/workflow/workflow_entry.py:163),它列出了跑一个工作流需要的全部"零件":
| 零件 | 是什么 |
|---|---|
graph_config | 原始 JSON({nodes, edges}),用户在画布上存的那份 |
graph: Graph | 把 JSON 解析成的结构化图对象(引擎真正跑的) |
variable_pool | ★变量池:所有节点的输入/输出都存这里,跨节点传数据靠它 |
graph_runtime_state | 运行时状态:变量池 + 开始时间 + 执行上下文,一次运行的"账本" |
command_channel | 指令通道:外部下"停止"命令走这里(默认内存通道) |
call_depth | 调用深度:工作流里可以嵌套调工作流,这个防无限套娃 |
variable_pool(变量池)★整个工作流的"共享内存"。节点 A 产出 {text: "..."} 存进池子,节点 B 用 {{#A.text#}} 从池子取。节点之间不直接说话,全靠变量池中转——这是理解工作流数据流的钥匙。graph vs graph_config一份是"图纸原文"(JSON),一份是"施工用的结构化模型"(Graph 对象)。引擎跑的是后者,但很多校验/调试要回看前者。graph_runtime_state一次运行的完整状态账本。为什么单独抽出来?因为工作流支持"暂停/恢复"(比如等人工审批),状态得能被完整保存和还原。call_depth 熔断构造函数第一件事就是 if call_depth > WORKFLOW_CALL_MAX_DEPTH: raise——工作流节点里能再调子工作流,不限深度会栈溢出。WorkflowEntry 就是车间主任,把这些凑齐、按下启动。WorkflowEntry:把引擎搭起来、挂上各种"层"
构造函数核心就是造一个 GraphEngine 并给它挂"层"(api/core/workflow/workflow_entry.py:213):
# api/core/workflow/workflow_entry.py:213
self.graph_engine = GraphEngine(
workflow_id=workflow_id,
graph=graph,
graph_runtime_state=graph_runtime_state,
command_channel=command_channel,
config=GraphEngineConfig( # 并发配置:几个 worker 并行跑节点
min_workers=dify_config.GRAPH_ENGINE_MIN_WORKERS,
max_workers=dify_config.GRAPH_ENGINE_MAX_WORKERS,
scale_up_threshold=dify_config.GRAPH_ENGINE_SCALE_UP_THRESHOLD,
scale_down_idle_time=dify_config.GRAPH_ENGINE_SCALE_DOWN_IDLE_TIME,
),
child_engine_builder=self._child_engine_builder, # 嵌套子工作流用
)
if dify_config.DEBUG: # ① 调试模式 → 挂"详细日志层"
self.graph_engine.layer(DebugLoggingLayer(...))
limits_layer = ExecutionLimitsLayer( # ② ★挂"执行限制层":熔断
max_steps=dify_config.WORKFLOW_MAX_EXECUTION_STEPS,
max_time=dify_config.WORKFLOW_MAX_EXECUTION_TIME)
self.graph_engine.layer(limits_layer)
self.graph_engine.layer(LLMQuotaLayer(tenant_id=tenant_id)) # ③ 挂"LLM 额度层"
if dify_config.ENABLE_OTEL or is_instrument_flag_enabled():
self.graph_engine.layer(ObservabilityLayer()) # ④ 挂"可观测层"(OTel)
GraphEngineConfig 并发引擎不是一个节点跑完再跑下一个——它用一个 worker 池 并行跑"互不依赖的节点",还能按负载自动扩缩容(scale up/down)。这就是为什么两条并行分支能同时跑。.layer(...) 层机制★关键设计:引擎核心只管跑图,"日志/限流/额度/观测"这些横切关注点用 层(Layer)挂上去,像给相机套滤镜。要加新能力,写个新 Layer 挂上即可,不动引擎核心。ExecutionLimitsLayer熔断层:限制"最多跑多少步""最多跑多久"。防止用户画了个死循环图把服务器跑挂。这是必挂的安全带。LLMQuotaLayer额度层:跑图过程中调 LLM 会消耗租户额度,这一层负责计量/拦截。把"额度"做成层,而不是塞进每个 LLM 节点,就是层机制的价值。run():为什么返回的是"事件流"
装配好后,执行就一句 run()(api/core/workflow/workflow_entry.py:250):
# api/core/workflow/workflow_entry.py:250
def run(self) -> Generator[GraphEngineEvent, None, None]: # ★返回类型是"事件生成器"
graph_engine = self.graph_engine
try:
# 在 Graphon 引擎的事件流上,套一层 Dify 的"响应流"兼容过滤
generator = iter_dify_graph_engine_events(graph_engine, self._response_stream_filter)
yield from generator # ① 引擎每产一个事件就 yield 出去
except GenerateTaskStoppedError:
pass # ② 用户主动停止 → 静默收尾
except Exception as e:
logger.exception("Unknown Error when workflow entry running")
yield GraphRunFailedEvent(error=str(e)) # ③ 出错也变成一个"失败事件"吐出
return
返回 Generator★注意返回类型 Generator[GraphEngineEvent]——不是"跑完给个结果",而是"边跑边一个个吐事件"。事件有很多种:节点开始、节点产出一块内容、节点结束、图跑完、图失败……前端订阅这个流,就能实时显示进度。iter_dify_graph_engine_events引擎(graphon)产的是"通用图事件",Dify 需要在上面套一层过滤,保住自己的"响应流语义"(比如哪些内容该实时流给用户、哪些不该)。这是 Dify 侧的适配层。GenerateTaskStoppedError用户点"停止"时,command_channel 传入停止命令,引擎抛这个异常。这里 pass 静默处理——停止是正常操作,不是错误。失败也是一个事件就算出未知异常,也不是直接崩,而是 yield GraphRunFailedEvent——把"失败"也统一成事件流里的一个事件。下游用同一套逻辑处理成功/失败,前端能优雅展示错误。拓扑:谁是起点、谁在谁的上游
图要跑,先得懂图的结构。graph_topology.py 把 JSON 解析成可查询的拓扑(api/core/workflow/graph_topology.py:16):
# api/core/workflow/graph_topology.py:16
class WorkflowGraphTopology:
@classmethod
def from_graph(cls, graph): # ① 从 {nodes, edges} 解析
node_ids = cls._node_ids_from_graph(graph) # 收集所有节点 id
edges = graph.get("edges")
for edge in edges: # 每条边: source → target
source = edge.get("source"); target = edge.get("target")
incoming[target].append(source) # 记录"target 的上游有谁"
return cls(node_ids=node_ids, incoming=incoming)
def is_upstream(self, *, source_node_id, target_node_id): # ② 判断 A 是不是 B 的上游
if source_node_id == target_node_id: return ...
queue = deque(self._incoming.get(target_node_id, ())) # 从 B 往上游 BFS
while queue:
candidate = queue.popleft()
if candidate == source_node_id: return True # 找到 A → 是上游
...
incoming 映射核心数据结构:{节点: [它的上游节点们]}。有了它,就能回答"这个节点依赖谁的输出""能不能开始跑(上游都跑完了吗)"。is_upstream (BFS)用广度优先从目标节点往上游走,判断某节点是否为其祖先。用途:校验变量引用合法性(你只能引用上游节点的输出,不能引用还没跑的下游)。upstream_node_ids(api/core/workflow/graph_topology.py:54)返回某节点所有上游节点集合。注意它 & self._node_ids 过滤——因为半删除的图里,边可能指向已不存在的节点,只返回真实存在的。起点怎么定拓扑本身不选起点;选起点是 Day09 的 get_default_root_node_id 干的(找类型是"开始节点"的那个)。拓扑负责"结构",起点是"语义"。upstream_node_ids 用 & self._node_ids 兜底,只认真实节点。处理用户数据时,永远别假设它是干净的——防御式解析是常态。single_step_run:只跑一个节点来调试
画布上你常需要"只测这一个节点对不对",不想跑整张图。WorkflowEntry.single_step_run(api/core/workflow/workflow_entry.py:265)就干这个:
# api/core/workflow/workflow_entry.py:265
def single_step_run(cls, *, workflow, node_id, user_id, user_inputs, variable_pool, ...):
node_config = workflow.get_node_config_by_id(node_id) # ① 只取这一个节点的配置
node_type = node_config["data"].type
node_version = str(node_config["data"].version)
node_cls = resolve_workflow_node_class( # ② 解析出它的节点类(Day09)
node_type=node_type, node_version=node_version)
run_context = build_dify_run_context( # ③ 造运行上下文(DEBUGGER 模式)
tenant_id=workflow.tenant_id, app_id=workflow.app_id,
user_id=user_id, invoke_from=InvokeFrom.DEBUGGER)
graph_runtime_state = GraphRuntimeState(variable_pool=variable_pool, start_at=..., ...)
if is_start_node_type(node_type): # ④ 若是开始节点,把用户输入灌进池子
add_node_inputs_to_pool(variable_pool, node_id=node_id, inputs=user_inputs)
...
GraphEngine、不挂那些层,只解析出一个节点类、造个最小运行时状态、把手填的 user_inputs 塞进变量池,然后单独跑这一个节点。变量池、运行时状态、节点解析这些零件是复用的,只是编排范围从"整张图"缩到"一个节点"。这让"调试单节点"和"运行整图"共享同一套底层机制——一致性带来可维护性。invoke_from=DEBUGGER 标记这是调试调用,日志/额度/追踪会据此区别对待。串起来 + 今日小结
[开始] → [知识检索] → [LLM] → [结束] 四个节点,三条边。→ 后端把这份 JSON 交给 WorkflowEntry(graph_config={nodes:4条, edges:3条}, graph=Graph对象, variable_pool=空池, runtime_state=新账本, call_depth=0) → 构造函数造 GraphEngine,挂上 ExecutionLimitsLayer(max_steps=..)、LLMQuotaLayer、(生产没开 DEBUG 就不挂日志层)→ 调 run():拓扑发现"开始"是根节点先跑(把用户问题存入变量池)→ "知识检索"从池取问题、检索、把结果存回池 → "LLM"从池取问题+检索结果拼 prompt(Day07!)、调 invoke_llm(Day05!)流式产出 → 每一步都 yield 一个事件 → task_pipeline 转 SSE → 前端看到"正在检索…正在生成…"。Day05/06/07 学的模型调用,正是这张图里"LLM 节点"内部干的事。👶 小白:GraphEngine 和 WorkflowEntry 到底谁跑图?
👨🏫 老师:分工明确。GraphEngine(在 graphon 包里)是真正的"发动机"——它懂图论:并行调度、按边推进、发事件、熔断。WorkflowEntry(在 Dify 的 core/workflow 里)是"装配 + 点火"——它把 Dify 特有的东西(变量池、各种 Layer、响应流过滤、租户额度)装配到通用引擎上,然后按 run()。为什么这么分?因为图执行是通用能力(可以独立成库 graphon),而"限流/额度/Dify 事件语义"是 Dify 业务。通用引擎 + 业务装配层,两边各自演进。Day10 我们会看引擎和 Dify 侧的运行时适配器怎么对接。
🧠 今天你应该能回答
- 跑一个工作流需要哪几样零件?(图/变量池/运行时状态/指令通道/调用深度)
- 节点之间怎么传数据?(全靠变量池中转,不直接说话)
- "层(Layer)"是干嘛的?(日志/限流/额度/观测等横切能力,可插拔挂在引擎外)
- run() 为什么返回 Generator?(事件流,边跑边播报,支撑实时进度)
- 拓扑能回答什么问题?(谁是谁的上游、能否引用某节点输出)
- single_step_run 和整图运行的关系?(复用同一套零件,编排范围缩到一个节点)
✋ 10 分钟动手
cd /Users/bitmart/work/codes/github/AI_WORK/dify
# 1. 装配入口 + 挂层
sed -n '163,248p' api/core/workflow/workflow_entry.py # __init__ 造引擎 + layer(...)
sed -n '250,262p' api/core/workflow/workflow_entry.py # run() 事件流
# 2. 拓扑解析
sed -n '16,83p' api/core/workflow/graph_topology.py # WorkflowGraphTopology
# 3. 单节点调试
sed -n '265,315p' api/core/workflow/workflow_entry.py # single_step_run
core/workflow/node_factory.py——register_nodes 怎么把内置节点(graphon.nodes)和 Dify 本地节点(core.workflow.nodes)注册到一起,resolve_workflow_node_class 怎么按"类型+版本"找到节点类,create_node 怎么按节点类型注入不同的依赖(LLM 节点要模型、工具节点要工具运行时)。