Day 08 / 共 20 天 · 阶段3 工作流

Workflow 总览:拖出来的一张流程图,是怎么被跑起来的

阶段2 讲清了"单次模型调用",但真实应用往往是一连串步骤:先检索知识库 → 再调 LLM → 再判断分支 → 再调工具。用户在画布上拖出这张"节点图",后端怎么执行?主角是 core/workflow/workflow_entry.py(入口)和 graph_topology.py(拓扑解析)。弄清三件事:①一个工作流由哪几样东西组成(图配置、Graph 对象、变量池、运行时状态);②WorkflowEntry 怎么把这些组装进 GraphEngine 并挂上各种"层"(限流/日志/观测);③执行为什么是"吐出一串事件"而不是"返回一个结果"。今天先建立整体世界观,Day09 钻节点、Day10 钻引擎。

📍 你在 20 天里的位置(阶段3:工作流 · D08-10)
D06 Provider/实体 D07 Prompt/生成 D08 工作流总览 D09 节点体系 D10 图执行引擎 S4 RAG S5 工具/Agent S6 收官
💡 先用两个类比兜住今天 类比一:一张工作流图像地铁线路图——每个"站"是一个节点(检索/LLM/判断/工具),"轨道"是边(谁连到谁),"哪站是起点"由拓扑决定。GraphEngine 就是那列按线路行驶的地铁,WorkflowEntry 是发车前把线路图、时刻表、乘客都装配好的调度室。类比二:执行工作流像看直播而不是看录播——引擎不是跑完给你一个大结果,而是"每到一站就播报一次"(节点开始了、节点产出了、出错了)。run() 返回的是一个源源不断的事件流(generator),前端才能实时看到"正在检索…正在生成…"。
L01

痛点:画布上拖出来的图,后端怎么执行

🤔 痛点用户在 Dify 画布上拖了 5 个节点、连了几条线,保存成一份 JSON({nodes:[...], edges:[...]})。后端拿到这份 JSON 后要解决一堆问题:从哪个节点开始跑?某节点的输入要用到前面节点的输出,数据怎么传递?多个节点能不能并行?跑到一半用户点了停止怎么办?执行进度怎么实时推给前端?跑太久/步数太多怎么熔断?如果这些逻辑和"具体某个节点干啥"混在一起,代码会彻底失控。
💡 本质:把"编排/调度"和"节点逻辑"彻底分开Dify 的答案是一套图执行框架(graphon.graph_engine.GraphEngine)+ Dify 侧的装配入口 WorkflowEntry。引擎只懂"图论":按边推进、并行调度、事件广播、熔断控制;至于"某个节点具体干什么",引擎一概不管,交给节点自己(Day09)。WorkflowEntry 的活是"装配"——把图、变量池、运行时状态、各种控制"层"拼到引擎上,然后按一个 run() 键。
L02

一张图由哪几样东西组成

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——工作流节点里能再调子工作流,不限深度会栈溢出。
大白话把工作流想成"一条流水线开工需要的东西":图纸(graph)、传送带(variable_pool,工件在工位间传递)、开工记录本(runtime_state)、急停按钮(command_channel)。WorkflowEntry 就是车间主任,把这些凑齐、按下启动。
L03

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 节点,就是层机制的价值。
💡 层(Layer)= 洋葱式的横切关注点日志、限流、额度、观测——这些能力和"业务逻辑"正交(哪个工作流都需要),但又不该硬编码进引擎。用"层"一圈圈包在引擎外面,可插拔、可组合、可按环境开关(比如 DEBUG 才挂日志层、开了 OTel 才挂观测层)。这和 Web 框架的"中间件"是同一个思想。
L04

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——把"失败"也统一成事件流里的一个事件。下游用同一套逻辑处理成功/失败,前端能优雅展示错误。
run() 吐出一串事件(直播,而非录播) WorkflowEntry.run() 节点开始 产出一块 节点结束 图跑完 task_pipeline 事件→SSE 前端实时进度 每到一站就播报 → 用户看到"正在检索…正在生成…"
图注:run() 是事件源头,一路 yield 到前端 SSE,接上了 Day04 的响应旅程。
L05

拓扑:谁是起点、谁在谁的上游

图要跑,先得懂图的结构。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_idsapi/core/workflow/graph_topology.py:54)返回某节点所有上游节点集合。注意它 & self._node_ids 过滤——因为半删除的图里,边可能指向已不存在的节点,只返回真实存在的。
起点怎么定拓扑本身不选起点;选起点是 Day09 的 get_default_root_node_id 干的(找类型是"开始节点"的那个)。拓扑负责"结构",起点是"语义"。
⚠️ 坑:边指向了不存在的节点源码注释专门提到 "Edges may reference ids missing from nodes (half-deleted graphs)"——用户删了个节点但连它的边没清干净,图数据就"半残"了。upstream_node_ids& self._node_ids 兜底,只认真实节点。处理用户数据时,永远别假设它是干净的——防御式解析是常态。
L06

single_step_run:只跑一个节点来调试

画布上你常需要"只测这一个节点对不对",不想跑整张图。WorkflowEntry.single_step_runapi/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)
    ...
💡 同一套零件,两种粒度对比 L03 的整图运行:single_step 不建整个 GraphEngine、不挂那些层,只解析出一个节点类、造个最小运行时状态、把手填的 user_inputs 塞进变量池,然后单独跑这一个节点。变量池、运行时状态、节点解析这些零件是复用的,只是编排范围从"整张图"缩到"一个节点"。这让"调试单节点"和"运行整图"共享同一套底层机制——一致性带来可维护性。
用途:画布上每个节点旁的"运行此步骤"按钮走的就是它。invoke_from=DEBUGGER 标记这是调试调用,日志/额度/追踪会据此区别对待。
L07

串起来 + 今日小结

📝 真实值:一张"客服知识库"工作流怎么跑 用户画了:[开始] → [知识检索] → [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
明日预告 · Day 09:今天引擎"跑节点"时,节点是从哪冒出来的?明天进 core/workflow/node_factory.py——register_nodes 怎么把内置节点(graphon.nodes)和 Dify 本地节点(core.workflow.nodes)注册到一起,resolve_workflow_node_class 怎么按"类型+版本"找到节点类,create_node 怎么按节点类型注入不同的依赖(LLM 节点要模型、工具节点要工具运行时)。
← Day 07 Prompt/生成 Day 09 · 节点体系 →