跳到主要内容

数据截至 (上游 commit b78a3462c9a6)

执行期横切:Layer 钩子、留痕落库与单节点调试

30 秒导读: 上一章(03)把 JSON 变成了可执行图,图交给 graphon 的 GraphEngine 去跑。这一章讲跑的时候旁边发生了什么:Dify 用 GraphEngineLayer 这个扩展点,把「落库、埋点、限时、写触发器日志」全部挂在引擎外面(模型额度已改到模型调用边界结算,见 §4.3);跑完在 Postgres 里留下五张表的记录;而画布上「单独跑一个节点」的调试体验,靠的是把每次运行的节点输出额外存成一份草稿变量


1. 这章解决的三个问题

先把问题摆清楚,后面每一节对应一个。

问题白话本章第几节
引擎跑的时候,谁在旁边听落库/埋点这些副作用挂在哪§3
跑到一半怎么让它停用户点「停止」,Web 进程怎么通知 Celery worker§4
跑完在数据库里剩下什么五张表 + 大字段怎么放§5
为什么能只跑一个节点上游变量从哪来§6

一句话先给结论:graphon 的 GraphEngine 只负责「按拓扑调度节点、发事件」,Dify 需要的一切外部副作用都做成 Layer 从外面挂上去。


2. 顶层全景:引擎、Layer、事件流

先看这张图。从左到右是数据流向;引擎在中间,上面挂 Layer(同步回调),右边流出事件给 SSE 管线(01 章讲过)。

┌──────── 挂在引擎上的 Layer(同步回调,跑在引擎线程里)────────┐
│ │
│ 落库 埋点(OTel) 限时/限步 │
│ Persistence Observability TimeSlice │
│ │ │ │ │
└─────┼──────────────┼──────────────┼────────────────────────┘
│ on_event / on_node_run_start / on_node_run_end

Graph ───────► GraphEngine(graphon,外部包) ───────► GraphEngineEvent 流
(03 章) ▲ │
│ 命令通道 CommandChannel ▼
│ (Pause / 停止) QueueManager → SSE

┌────────┴─────────┐
│ │
RedisChannel InMemoryChannel
(跨进程停止) (调试单跑缺省)

怎么读这张图: 事件是单向广播(引擎 → Layer 和 SSE),命令是反向的一条细线(外部 → 引擎)。Layer 既能收事件,也能往命令通道发命令——TimeSliceLayer 就是这么让引擎暂停的。

一个必须先说清的边界: GraphEngineGraphEngineLayerRedisChannelInMemoryChannel 都不在 Dify 仓库里,它们来自外部依赖 graphon==0.7.0api/pyproject.toml:48)。所以本章能引的是 Dify 侧的实现和用法;graphon 内部怎么调度线程、怎么消费命令,克隆里看不到,我不编。


3. GraphEngineLayer:扩展点长什么样

3.1 它要解决的小问题

工作流引擎的核心逻辑是「按依赖顺序跑节点」。但真实产品里,每跑一个节点都要顺带干一堆事:写一行 workflow_node_executions、开一个 OTel span、看看是不是跑太久该暂停了。

如果这些都写进引擎,引擎就和 Dify 的数据库、Redis、计费系统绑死了。Layer 就是那条切口:引擎发事件,副作用订阅事件

3.2 五个钩子

Dify 的各个 Layer 一共重写了这五个方法(从 Dify 侧的 @override 反推,graphon 基类本身不在克隆里):

钩子什么时候被调谁用了它
on_graph_start()整张图开跑前全部 Layer(一般拿来清空内部缓存)
on_event(event)每个 GraphEngineEvent 到达Persistence / TriggerPost / ConversationVariablePersistence
on_node_run_end(node, error, result_event)单个节点跑完Observability(结束 span)
on_node_run_start(node)单个节点开跑前Observability(开 span)
on_graph_end(error)整张图结束TimeSlice(撤销定时任务)、Observability(查漏未关闭的 span)

Layer 拿到的上下文靠一次 initialize 注入——单测里能看到这个签名:layer.initialize(read_only_state, command_channel=None)api/tests/unit_tests/core/app/workflow/test_persistence_layer.py:111)。注入之后,Layer 里就能用两个属性:

  • self.graph_runtime_state —— 只读的运行时状态(变量池、token 计数、步数)。TriggerPostLayer 就是从这里取 total_tokensoutputs 的(api/core/app/layers/trigger_post_layer.py:75-79on_event)。
  • self.command_channel —— 反向命令通道,见 §4。

3.3 Dify 的 Layer 全家福

Dify 侧一共五个业务 Layer,分两个目录放,划分依据是「依赖多重」:

Layer文件(api/ 下)干什么挂在哪
WorkflowPersistenceLayercore/app/workflow/layers/persistence.py:83把运行/节点执行写进数据库,顺带投递 trace 任务app_runner 里手动挂
ObservabilityLayercore/app/workflow/layers/observability.py:43每个节点开一个 OpenTelemetry spanWorkflowEntry.__init__ 按开关挂
TimeSliceLayercore/app/layers/timeslice_layer.py:16定时检查是否超配额,超了就发 PAUSE异步触发任务传入
ConversationVariablePersistenceLayercore/app/layers/conversation_variable_persist_layer.py:24conversation.* 变量更新落库只有 advanced-chat 挂
TriggerPostLayercore/app/layers/trigger_post_layer.py:24终态时回写 workflow_trigger_logs异步触发任务传入

(另有 PauseStatePersistenceLayercore/app/layers/pause_state_persist_layer.py:77)与 SuspendLayercore/app/layers/suspend_layer.py:7)属于暂停恢复,见 05 章。)

挂载的两条路,看 api/core/app/apps/workflow/app_runner.py:200-214

workflow_entry.graph_engine.layer(persistence_layer) # 运行器写死的
workflow_entry.graph_engine.layer(build_workflow_agent_workspace_retirement_layer(...))
for layer in self._graph_engine_layers: # 调用方传进来的
workflow_entry.graph_engine.layer(layer)

也就是说:落库是每次运行都要的,写死;限时/触发器日志是特定入口才要的,由调用方注入。 异步触发的注入点在 api/tasks/async_workflow_tasks.py:283-284TimeSliceLayer + TriggerPostLayer)。

WorkflowEntry.__init__ 自己会挂引擎级 Layer(api/core/workflow/workflow_entry.py:156-175):开 DEBUG 时挂 DebugLoggingLayer(来自 graphon)、无条件挂 ExecutionLimitsLayer(限步数/限时长,同样来自 graphon)、开了 OTel 才挂 ObservabilityLayer。模型额度已不在此列——它从 Layer 体系搬到了模型调用边界(§4.3)。

3.4 落库 Layer 怎么工作

WorkflowPersistenceLayer 是这套设计最典型的样本。它三个生命周期钩子的分工非常干脆:

  • on_graph_start()persistence.py:111)—— 只做清空:清掉 _node_execution_cache_node_snapshots、序号计数器。一行数据库都不写。
  • on_event(event)persistence.py:118)—— 一个大 match,十几种事件各自对应一个 _handle_*。所有落库都在这里发生。
  • on_graph_end(error)persistence.py:146)—— 直接 return,什么都不做。

第三条值得停一下:图结束时它不写任何东西。因为终态早就由 GraphRunSucceededEvent / GraphRunFailedEvent / GraphRunAbortedEvent 这些事件在 on_event 里处理完了。事件是唯一的事实来源,生命周期钩子只管内存状态。

节点级的两个典型:

  • _handle_node_startedpersistence.py:226):构造一个 WorkflowNodeExecution 领域对象(状态 RUNNING),存进内存缓存,repository.save(...) 落一条「跑起来了」的记录,再往前端的 inspector 频道推一条 status="running"
  • _update_node_executionpersistence.py:424):节点成功/失败/异常都走这里,算 elapsed_time、写 outputs,然后连着调两个仓储方法save()save_execution_data()。为什么要分两次,§5.3 讲。

3.5 这个设计的好处与代价

好处比较明显:

  • 引擎可复用。 graphon 是独立包,不知道 Postgres、不知道租户、不知道 Redis。Dify 自己的 single_step_run 调试路径连一个 Layer 都不挂也能跑。
  • 副作用可裁剪。 调试运行不写 workflow_app_logs,没开 OTel 就不挂 ObservabilityLayerapi/core/workflow/workflow_entry.py:173-175 的开关判断),都是删一行的事。
  • 测试便宜。 单测直接 layer.initialize(fake_state, command_channel=None) 然后喂事件,不用起引擎。

代价也是真的:

  1. Layer 跑在引擎线程里。 TriggerPostLayer.on_event 里直接开了 DB session 并 commit()trigger_post_layer.py:53-91)——这是同步阻塞引擎的。Dify 对此的应对是在把执行搬进新线程之前先 db.session.close() 归还连接,注释写得很直白:「新线程里的操作可能跑很久」(api/core/app/apps/workflow/app_generator.py:379-382)。
  2. 钩子不够用时只能捅私有方法——Dify 的解法是把这类职责搬出 Layer。 旧版 LLMQuotaLayer 要在节点跑之前就让它失败,但 graphon 没给「跑前失败/跳过」的公开钩子,于是它直接把 node._run 换成一个返回 FAILED 的闭包并留 TODO 等上游开放公开钩子。上游后来把这一层整个移除,配额改为在模型调用边界结算(§4.3)——靠私有 override 顶住的需求,最终以「换挂载点」而不是「等钩子」收场。这是「扩展点边界画得不够宽」时最常见的症状与出路。
  3. 顺序耦合。 Layer 按 layer() 的调用顺序挂载,落库 Layer 永远第一个挂。代码里没有优先级声明,谁先谁后取决于挂载顺序 (inferred)。

4. 运行中控制:命令通道的两种接法

4.1 它要解决的小问题

用户在浏览器点「停止」。这个 HTTP 请求打到某个 Web 进程;而工作流可能跑在另一台机器的 Celery worker 里。两个进程之间怎么传一句「停」?

4.2 两种通道,两种场景

跨进程(生产运行) 进程内(调试 / 子图)
───────────────────── ────────────────────
[Web 进程] [同一个进程]
POST .../stop WorkflowEntry.__init__
│ │ 未传 command_channel
▼ ▼
GraphEngineManager(redis) InMemoryChannel()
.send_stop_command(task_id) │
│ 写 Redis key │ 内存队列
▼ "workflow:{task_id}:commands" ▼
[Worker 进程] RedisChannel 轮询 ──► GraphEngine ◄── 直接读
通道构造位置key / 载体适用
RedisChannelapi/core/app/apps/workflow/app_runner.py:164api/core/app/apps/advanced_chat/app_runner.py:234f"workflow:{task_id}:commands"生产运行,Web 进程和执行进程可能不同
InMemoryChannelapi/core/workflow/workflow_entry.py:135(缺省)进程内对象单节点调试、以及未显式传通道的一切入口

两个 app_runner 都以 Redis 通道为底,但写法已经分化:workflow 侧把 key 生成收进 app_task_command_channel_keyapi/core/app/apps/workflow/app_runner.py:156-164),advanced-chat 侧则把 RedisChannel 和一个监听 Celery warm-shutdown 信号的 CelerySignalCommandChannel 组合成 CombinedCommandChannelapi/core/app/apps/advanced_chat/app_runner.py:226-238)——worker 收到优雅停机信号时也能像用户点停止一样中断图。

# workflow 侧(示意)
command_channel = RedisChannel(redis_client, app_task_command_channel_key(task_id))
# advanced-chat 侧(示意):Redis 命令 + Celery 停机信号,两路合一
command_channel = CombinedCommandChannel((RedisChannel(redis_client, channel_key), celery_signal_channel))

channel key 只由 task_id 决定——这就是为什么前端只要拿着 task_id 就能停掉一次运行,不需要知道它跑在哪台机器上。

发送侧散落在各个入口控制器里,写法统一(例如 api/controllers/console/app/workflow.py:1201api/controllers/web/workflow.py:131api/services/app_task_service.py:46):

GraphEngineManager(redis_client).send_stop_command(task_id)

这里只是「停止」的一半。 同一个停止接口还会写一个 Redis 标志位 generate_task_stopped:{task_id}api/core/app/apps/base_app_queue_manager.py:233),由 AppQueueManager.listen() 每秒轮询时发现并补发停止事件——那是队列侧的旧机制,本节讲的命令通道是引擎侧的新机制,两者为向后兼容并存。队列侧那一半的细节见 01 章 §7.2

WorkflowEntry 的缺省行为是「没给通道就自己造一个内存的」(api/core/workflow/workflow_entry.py:134-135),所以调试路径完全不需要 Redis 也能跑。迭代/循环子图不再由 Dify 侧单独建子引擎——_WorkflowChildEngineBuilder 已移除,容器(迭代/循环)拓扑直接进同一张图、由同一个引擎调度(api/core/workflow/generator/runner.py:1280 起的容器合成逻辑),天然同进程。

4.3 谁在往通道里发命令

除了用户点停止,Layer 自己也是命令的发送方。这是 Layer 设计里最巧的一环:它既是观察者,又能反向干预。

发送方命令触发条件代码
停止 APIstop用户点击controllers/console/app/workflow.py:1201
TimeSliceLayerCommandType.PAUSEAPScheduler 定时器发现资源配额到顶core/app/layers/timeslice_layer.py:54

TimeSliceLayer 的实现挺有意思:它在 on_graph_start 里往一个类级别共享的 BackgroundScheduler 注册一个周期任务(timeslice_layer.py:67-80),周期是 plan.granularity 秒;任务每次醒来问一句 cfs_plan_scheduler.can_schedule(),返回 RESOURCE_LIMIT_REACHED 就发 PAUSE 并把自己从调度器摘掉。on_graph_end 负责兜底删任务(timeslice_layer.py:87-91)。

模型额度的新家更值得记:预留-提交-释放三段式结算,全部收在模型调用边界。LLMQuotaLayer 已移除,配额改由 QuotaManagedModelInstanceapi/core/model_manager.py:446)负责——invoke_llm_reserve_quota_for_request 预留,拿到响应后按 usage 提交,finallyrelease_quota_safely 释放(model_manager.py:551-567);流式调用分「边收边提交」与「攒完再交」两种模式(_invoke_llm_streammodel_manager.py:569 起)。轮询式 LLM 的结算在 workflow 侧的 DifyPreparedPollingLLM._settle_polling_quotaapi/core/workflow/node_runtime.py:332)。额度不足时 reserve_model_quota_for_model 直接抛 QuotaExceededErrorapi/core/app/llm/quota.py:139)——失败落在这一次模型调用上,而不是像旧 Layer 那样反向发 AbortCommand 中止整张图。


5. 留痕:跑完在数据库里剩下什么

5.1 五张表,各管一段

表(模型)api/models/workflow.py一行代表谁写的
workflow_runsWorkflowRun:743一次完整运行:状态、耗时、token、输入输出WorkflowPersistenceLayer 经仓储
workflow_node_executionsWorkflowNodeExecutionModel:922一个节点的一次执行同上
workflow_node_execution_offloadWorkflowNodeExecutionOffload:1171某次执行被卸载到对象存储的大字段仓储在 save_execution_data
workflow_app_logsWorkflowAppLog:1280面向「应用日志」列表的一条记录(不含调试SSE 管线,generate_task_pipeline.py:766
workflow_archive_logsWorkflowArchiveLog:1372归档后的运行快照(run + log + trigger 三方字段拍平)保留期任务

前两张是执行期实时写的;后两张是「事后」的。

为什么 workflow_app_logs 要独立于 workflow_runs 因为它只收「真用户跑的」:generate_task_pipeline.py:750-761 里那个 match invoke_from,遇到 DEBUGGER / TRIGGER / PUBLISHED_PIPELINE / VALIDATION 直接 return 不写。调试运行照样进 workflow_runstriggered_from=DEBUGGING,见 api/models/enums.py:24-31),但不会污染应用日志列表。

WorkflowArchiveLog 则是反范式的快照:它把 WorkflowRun 的字段加 run_ 前缀(run_status / run_elapsed_time / run_total_tokens…)、把 WorkflowAppLog 的字段加 log_ 前缀,再塞一个 trigger_metadatamodels/workflow.py:1453-1474)。这样原始运行记录被清理后,日志列表还能显示。

5.2 一次成功运行的写入时序

GraphRunStartedEvent ──► WorkflowExecution.new(...) ──► workflow_runs INSERT (running)
NodeRunStartedEvent ──► WorkflowNodeExecution(RUNNING) ──► node_executions INSERT (running)
NodeRunSucceededEvent ──► save() + save_execution_data() ──► node_executions UPDATE + 大字段卸载
…每个节点重复…
GraphRunSucceededEvent ──► 汇总 token/steps/outputs ──► workflow_runs UPDATE (succeeded)
└─► TraceQueueManager.add_trace_task(...) (§7)
SSE 管线收尾 ──► WorkflowAppLog(...) ──► workflow_app_logs INSERT

汇总数据不是 Layer 自己数的,而是从只读运行时状态里抄的(persistence.py:415-421_populate_completion_statistics):runtime_state.total_tokensruntime_state.node_run_stepsruntime_state.exceptions_count

5.3 两段式写入:savesave_execution_data

仓储接口上有两个方法(api/core/repositories/factory.py:35-40):

class WorkflowNodeExecutionRepository(Protocol):
def save(self, execution: WorkflowNodeExecution): ...
def save_execution_data(self, execution: WorkflowNodeExecution): ...

save() 写元数据(状态、时间、序号);save_execution_data() 才处理 inputs / outputs / process_data 这些可能巨大的字段。终态时两个连着调(persistence.py:462-463)。

拆开的原因写在 save() 的注释里(api/core/repositories/sqlalchemy_workflow_node_execution_repository.py:337-340):引擎对同一个节点会多次调 save——开跑一次、每次重试一次、终态再一次——只有最后一次带全 inputs/outputs,前面几次必须容忍缺数据、不能尝试卸载。

5.4 大字段卸载(offload)

问题: 一个 LLM 节点的 outputs 可能是几 MB 的文本,一个知识检索节点的 inputs 可能是上千条 chunk。全塞进 Postgres 的 LongText,查列表页都会被拖死。

做法: 超过阈值就把完整值写进对象存储,数据库里只留截断后的值 + 一条指针记录。

outputs(原始)


truncate_variable_mapping() ──► 没超阈值 ──► 直接 JSON 写进 node_executions.outputs
│ 超了

① 完整 JSON 上传 → UploadFile
② node_executions.outputs = 截断值 (列表页/前端读这个)
③ 插一行 workflow_node_execution_offload(type=outputs, file_id=…)

实现在 _truncate_and_uploadsqlalchemy_workflow_node_execution_repository.py:277)和 save_execution_data(同文件 :405),三种类型各走一遍:INPUTS / OUTPUTS / PROCESS_DATA(枚举见 api/models/enums.py:51-54)。阈值是 WORKFLOW_VARIABLE_TRUNCATION_MAX_SIZE,默认 1000 KiB(api/configs/feature/__init__.py:861-865)。

读回来时,模型上有一组对称的方法:inputs_truncated / outputs_truncated 判断是否被截断,load_full_inputs(session, storage) / load_full_outputs(...) 按需从对象存储拉全量(models/workflow.py:1174-1217)。

两个设计细节值得记:

  • inputs 和 outputs 分开存,不合并成一个对象。 模型文件里有一大段注释解释为什么(models/workflow.py:1238-1259):合并需要缓冲第一次 save 到执行结束才 flush,那样节点执行状态在完成前就不可观测了——「显著损害可观测性」,所以宁可多一次 I/O。
  • node_execution_id 可以为 NULL,表示这条卸载记录已经和执行记录脱钩,等垃圾回收(models/workflow.py:1240-1243)。唯一约束靠 PostgreSQL「NULL 互不相等」的默认行为,才允许多条 NULL 并存(同文件 :1174-1186 的注释)。

5.5 仓储抽象:两层,不是一层

repositories/ 目录下东西不少,但结构其实很整齐——按「谁用」分成两层

core/repositories/ 写入侧(引擎在跑时用)
factory.py Protocol 定义 + DifyCoreRepositoryFactory
sqlalchemy_workflow_execution_repository.py 同步写
sqlalchemy_workflow_node_execution_repository.py 同步写 + offload
celery_*_repository.py 异步写(丢给 Celery 任务)

│ 继承
repositories/ 读取侧(Service / Controller 用)
factory.py DifyAPIRepositoryFactory
api_workflow_run_repository.py Protocol:分页/统计/清理/归档
sqlalchemy_api_workflow_run_repository.py 实现
api_workflow_node_execution_repository.py Protocol
sqlalchemy_api_workflow_node_execution_repository.py
execution_extra_content_repository.py Protocol(按 message_id 取附加内容)
工厂接口形态典型方法
写入侧DifyCoreRepositoryFactorycore/repositories/factory.py:55极窄:save / save_execution_data / get_by_workflow_execution引擎线程里调
读取侧DifyAPIRepositoryFactoryrepositories/factory.py:17,继承前者)很宽:分页、按时间批量取、软删、归档、暂停记录get_paginated_workflow_runscreate_archive_logsget_expired_runs_batch

实现类是配置字符串,不是硬编码。 工厂用 import_string(class_path) 动态加载(core/repositories/factory.py:86-94),路径来自四个配置项(api/configs/feature/__init__.py:955-980):

配置项默认实现换成什么
CORE_WORKFLOW_EXECUTION_REPOSITORYSQLAlchemyWorkflowExecutionRepositoryCeleryWorkflowExecutionRepository
CORE_WORKFLOW_NODE_EXECUTION_REPOSITORYSQLAlchemyWorkflowNodeExecutionRepositoryCeleryWorkflowNodeExecutionRepository
API_WORKFLOW_NODE_EXECUTION_REPOSITORYDifyAPISQLAlchemyWorkflowNodeExecutionRepository自定义
API_WORKFLOW_RUN_REPOSITORYDifyAPISQLAlchemyWorkflowRunRepository自定义

Celery 版的意义:把落库这个阻塞操作从引擎线程挪到后台 worker,代价是「刚写的立刻读」需要一层内存缓存兜着(api/core/repositories/celery_workflow_node_execution_repository.py:38-44 的类注释)。这也是 §3.5 那条「Layer 阻塞引擎线程」的官方解药。

ExecutionExtraContentRepositoryrepositories/execution_extra_content_repository.py:9)是个只有一个方法的极小 Protocol,按 message_ids 批量取附加内容(实现里主要处理人工输入表单,见 05 章)。


6. 调试体验:单节点为什么能独立跑

6.1 它要解决的小问题

画布上一个 LLM 节点,输入引用了上游 Start 节点的 {{#start.query#}}。你点它右上角的「运行此步骤」——上游根本没跑,那个变量的值从哪来?

答案:从上一次跑留下的草稿变量里来。

6.2 草稿变量:调试态的变量快照

调试模式跑一次完整工作流

│ 每个节点成功后

DraftVariableSaver.save(process_data, outputs)


workflow_draft_variables 表 (超大值 → workflow_draft_variable_files → 对象存储)
唯一键 (app_id, user_id, node_id, name)

│ 下次单节点调试时

DraftVarLoader.load_variables(selectors) ──► 灌进 VariablePool ──► 节点可以独立跑

两张表:

位置存什么
workflow_draft_variablesWorkflowDraftVariablemodels/workflow.py:1559一个变量的当前草稿值
workflow_draft_variable_filesWorkflowDraftVariableFilemodels/workflow.py:2029被卸载的大变量的元数据(size / length / 原始 value_type)

WorkflowDraftVariable 上有几个非常「产品化」的字段:

  • node_id 是个复用字段:普通节点填节点 id,会话变量填 conversation,系统变量填 sysmodels/workflow.py:1624-1629 的注释)。
  • visible 决定要不要在变量检查面板里显示,editable 决定用户能不能改(:1573-1579)。IF_ELSE 节点的输出一律不可见、不可编辑的系统变量一律不可见(services/workflow_draft_variable_service.py:1191-1196_should_variable_be_visible)。
  • last_edited_atNone 表示「创建后没被人改过」(:1529-1535)。
  • file_id 非空表示这个值被卸载了,value 里是截断版:1567-1571)。

构造必须走三个工厂方法 new_conversation_variable / new_sys_variable / new_node_variablemodels/workflow.py:1869/1891/1913),类文档明确禁止直接用构造器——因为要维护一堆不变式。

6.3 保存侧:Protocol + Factory + Noop

保存这件事被切成了「端口 / 适配器」:

角色位置说明
DraftVariableSaver(Protocol)core/app/apps/draft_variable_saver.py:10只有一个 save(process_data, outputs)
DraftVariableSaverFactory(Protocol)core/app/apps/draft_variable_saver.py:17(app_id, node_id, node_type, node_execution_id, enclosing_node_id) 造一个 saver
NoopDraftVariableSavercore/app/apps/draft_variable_saver.py:31什么都不做
_DebuggerDraftVariableSavercore/app/apps/base_app_generator.py:35开 Session,转调真实实现
DraftVariableSaver(真实实现)services/workflow_draft_variable_service.py:827建变量对象、批量 upsert

分派逻辑就一句话(core/app/apps/base_app_generator.py:332-360_get_draft_var_saver_factory):

if invoke_from == InvokeFrom.DEBUGGER:
# 造 _DebuggerDraftVariableSaver
else:
# 造 NoopDraftVariableSaver

非调试运行一个草稿变量都不写。 这是空对象模式(Null Object)的教科书用法——调用方(SSE 管线)不需要判断模式,无脑调 saver.save(...) 就行(core/app/apps/workflow/generate_task_pipeline.py:793-800_save_output_for_event)。

真实实现 DraftVariableSaver.save()services/workflow_draft_variable_service.py:1162)按节点类型分三路:

节点类型数据来源方法
VARIABLE_ASSIGNERprocess_data_build_from_variable_assigner_mapping
START / 触发器类outputs(要做名字规范化)_build_variables_from_start_mapping
其它outputs_build_variables_from_mapping

最后统一 _batch_upsert_draft_variable:662)按唯一键覆盖。

两条不保存的规则,外加一条反向的兜底:

  • 迭代/循环内部的节点不保存(除了变量赋值器)——_should_save_output_variables_for_draft:906-911)。所以 _enclosing_node_id 这个参数不是装饰。
  • 部分变量按节点类型排除:LLM 的 finish_reason、Loop 的 loop_round:833-840_EXCLUDE_VARIABLE_NAMES_MAPPING,注释在 :833-836、定义在 :837-840)。
  • 反过来,节点没有任何输出时,会塞一个 __dummy__ 的不可见变量当「我跑过了」的信号(:830-831:894-905)。

6.4 加载侧:DraftVarLoader

DraftVarLoaderservices/workflow_draft_variable_service.py:80)实现 graphon 的 VariableLoader 接口,load_variables(selectors) 按选择器批量查草稿变量。里面两个细节:

  • File 类型要二次加载。 文件段(FileSegment / ArrayFileSegment)里的 storage_key 不在草稿变量里,得用 StorageKeyLoader 再查一遍(:124-134)。
  • 被卸载的变量用线程池并发拉。 ThreadPoolExecutor(max_workers=10) 并发调 _load_offloaded_variable:158-163),因为每个都要打一次对象存储。

6.5 单节点跑:从 HTTP 到 single_step_run

POST /apps/{app_id}/workflows/draft/nodes/{node_id}/run
controllers/console/app/workflow.py:1028 DraftWorkflowNodeRunApi.post


services/workflow_service.py:867 run_draft_workflow_node
│ ① 预填会话变量默认值
│ ② 造 VariablePool(Start 类节点还要建/取 conversation)
│ ③ 造 DraftVarLoader
│ ④ 算 enclosing_node_id(节点是否在迭代/循环里)

core/workflow/workflow_entry.py:253 WorkflowEntry.single_step_run
│ ⑤ 解析节点类、算变量映射
│ ⑥ load_into_variable_pool(...) ← 缺的变量从草稿里补
│ ⑦ DifyNodeFactory 造出单个 Node,直接跑
▼ (注意:没有 Graph、没有 GraphEngine、没有 Layer)
回到 workflow_service.py
│ ⑧ repository.save(node_execution) triggered_from=SINGLE_STEP
│ ⑨ DraftVariableSaver.save(...) 把这次的输出也存成草稿变量

返回 WorkflowNodeExecutionModel

关键点:single_step_run 根本不走引擎。 它用 DifyNodeFactory 造出一个 Node 就直接调,返回 (node, generator)workflow_entry.py:289-302)。所以:

  • 一个 Layer 都不挂,落库是 workflow_service 手动做的(services/workflow_service.py:1229-1236)。
  • workflow_node_executions.workflow_run_id 为 NULL——模型注释写明「单步调试时为空」(models/workflow.py:986-988);triggered_fromSINGLE_STEPmodels/workflow.py:965)。
  • 跑完还要再存一次草稿变量(services/workflow_service.py:1246-1256),这样下游节点下次单跑时就能引用到本次的输出。闭环就是这么合上的。

6.6 单次迭代 / 单次循环:这个走引擎

「单独跑一次迭代」和单节点不一样——迭代体里可能有好几个节点,必须真的调度。所以它走的是完整的 generator 路径,只是把图裁小了

入口:WorkflowAppGenerator.single_iteration_generatecore/app/apps/workflow/app_generator.py:431)和 single_loop_generate:495),advanced-chat 有对应的一对(core/app/apps/advanced_chat/app_generator.py:320 / :391)。它们做的事:

  1. 造一个 invoke_from=DEBUGGERWorkflowAppGenerateEntity,带上 SingleIterationRunEntity(node_id, inputs)app_generator.py:447-449)。
  2. 仓储用 WorkflowRunTriggeredFrom.DEBUGGING + WorkflowNodeExecutionTriggeredFrom.SINGLE_STEP:461-472)。
  3. DraftVarLoader,走正常的 _generate

裁图发生在 runner 里(core/app/apps/workflow_app_runner.py:240_get_graph_and_variable_pool_for_single_node_run)。逻辑很朴素:

# 示意,非源码:只留下「迭代节点自己 + 属于它的子节点 + 它的起始节点」
node_configs = [
node for node in graph_config["nodes"]
if node["id"] == node_id # 迭代节点本身
or node["data"].get("iteration_id", "") == node_id # 挂在它下面的子节点
or node["id"] == start_node_id # 迭代体的入口
]
# 边同理:两端都必须在保留的节点集合里

node_type_filter_key 参数就是 "iteration_id""loop_id" 的开关(workflow_app_runner.py:216 / :223)。裁完的 graph_config 交给 Graph.init 正常建图,后面和普通运行完全一样——Layer 照挂,SSE 照流

对照记一下三种调试的差别:

方式走引擎吗挂 Layer 吗workflow_run_id图的范围
单节点(single_step_runNULL只有那个 Node 对象
单次迭代 / 单次循环裁剪后的子图
完整调试运行全图

7. 可观测:两条互不相干的 trace 通路

新手最容易搞混的地方:Dify 里「trace」有两套,目标不同、路径不同、开关不同

ObservabilityLayerTraceQueueManager
面向谁运维(APM)应用开发者(LLM 可观测平台)
协议OpenTelemetryLangfuse / LangSmith 等各家 SDK
粒度每个节点一个 span每次运行一个 trace task
触发点on_node_run_start / on_node_run_end图终态时 _enqueue_trace_task
同步性同步,进程内异步,落盘 + Celery
开关dify_config.ENABLE_OTEL 或 instrument flag应用配了 ops trace provider

7.1 ObservabilityLayer:给每个节点开一个 span

on_node_run_startcore/app/workflow/layers/observability.py:92)干三件事:用节点标题开 span、context_api.attach(new_context) 把 span 塞进当前 OTel 上下文、记下 (span, token)

中间那步是精华。 一旦 span 进了上下文,节点里发出的所有 HTTP 请求、数据库查询就会被 OTel 的自动埋点自动挂到这个节点的 span 下面——不需要在 HTTP 节点、LLM 节点里写任何埋点代码。文件头注释把这个意图讲得很直接(observability.py:3-6)。

on_node_run_end:124)按节点类型选解析器写属性再关 span。解析器注册表只特化了三类(:74-80):TOOL / LLM / KNOWLEDGE_RETRIEVAL,其余走 DefaultNodeOTelParser

两个防御细节:_init_tracer 在构造器里就试着拿 tracer,拿不到就把 _is_disabled 置真,之后所有钩子直接 return(:60-70)——关掉 OTel 时开销接近零on_graph_end 会检查还有没有没关掉的 span,有就打 warning(:166-173)。

7.2 TraceQueueManager:批量攒、定时刷、丢给 Celery

它和 Layer 的关系是单向的WorkflowPersistenceLayer 构造时可以接一个 trace_managerpersistence.py:92),图跑到终态时调 _enqueue_trace_taskpersistence.py:476)造一个 TraceTask(TraceTaskName.WORKFLOW_TRACE, ...) 丢进去。

TraceQueueManagercore/ops/ops_trace_manager.py:1512)内部是「攒批 + 定时器」:

add_trace_task() ──► 模块级全局 queue.Queue

threading.Timer(默认 5 秒,TRACE_QUEUE_MANAGER_INTERVAL)
│ 到点

collect_tasks() 最多取 100 条(TRACE_QUEUE_MANAGER_BATCH_SIZE)


send_to_celery()
│ ① task.execute() 算出 trace_info
│ ② 序列化后 storage.save(ops_trace/{app_id}/{uuid}.json)
│ ③ process_trace_tasks.delay({file_id, app_id})

Celery worker 真正发给第三方

(队列/定时器/批量常量在 ops_trace_manager.py:1506-1509send_to_celery:1542。)

为什么要先落盘再发 Celery? 因为 trace payload 可能很大,直接塞进 Celery 消息体不合适——所以走「存储放大件、消息传小指针」这个经典套路(:1554-1568)。

add_trace_task 还有个短路:没配 trace 实例、也没开企业遥测,就直接不入队(:1508)——没配置的用户完全不付出代价

7.3 第三条:给前端的 inspector 频道

除了上面两条,落库 Layer 还会往一个 Redis pub/sub 频道推节点状态变化:_inspector_publish_node_changed(workflow_run_id, node_id, status)persistence.py:264:276 等多处,实现在 api/services/workflow/inspector_events.py:134)。这条是给前端变量检查面板用的,和 trace 无关。


8. 巧妙之处(可以直接借鉴的)

  1. 生命周期钩子只管内存,事件才是事实来源。 WorkflowPersistenceLayer.on_graph_start 只清缓存、on_graph_end 直接 returnpersistence.py:111:143)。所有落库都由事件驱动,于是「暂停后恢复」这种半程场景不需要为钩子写特例。

  2. 空对象消灭调用点的 if。 NoopDraftVariableSavercore/app/apps/draft_variable_saver.py:31)让 SSE 管线无脑调 saver.save(...),是否调试模式的判断被收敛到工厂一处(base_app_generator.py:332)。

  3. 可观测性优先于 I/O 效率的显式取舍。 inputs/outputs 本可以合并成一次卸载,Dify 选择分开,理由是合并会让节点在完成前不可观测——注释写了 20 行来解释(models/workflow.py:1238-1259)。把「为什么没做那个显然的优化」写下来,比优化本身更有价值。

  4. Layer 既是观察者又是干预者。 TimeSliceLayer 通过 self.command_channel 反向发 PAUSE(timeslice_layer.py:54)。扩展点给了双向能力,避免了「为了停机再开一个后门」。

  5. channel key 只由 task_id 决定。 f"workflow:{task_id}:commands"app_runner.py:149)—— 停止请求不需要知道运行在哪台机器上,天然支持水平扩容。

  6. 仓储实现是配置字符串。 同一份 Layer 代码,改一个环境变量就从同步落库切成 Celery 异步落库(configs/feature/__init__.py:955-969)。这是解决「Layer 阻塞引擎线程」的现成开关。


9. 边界与坑

  • Layer 同步跑在引擎线程里。 TriggerPostLayer.on_event 里开 session 并 commit(trigger_post_layer.py:53-91),落库 Layer 每个节点两次写库。图越大,引擎线程被 I/O 拖住的时间越多。缓解手段是 Celery 仓储和「进新线程前先关掉 Flask session」(app_generator.py:379-382)。

  • TimeSliceLayer 用了类级别共享的 APScheduler。 scheduler: ClassVar[BackgroundScheduler]timeslice_layer.py:21),同进程内所有工作流共用一个后台调度器。任务 id 是随机 hex,on_graph_end 负责删;如果 on_graph_end 没被调到(进程崩了),任务会残留 (inferred)。而且它当前在同步触发路径上是被注释掉的——api/tasks/async_workflow_tasks.py:182 写着「TODO: Re-enable TimeSliceLayer after the HITL release」。

  • 单节点调试和真实运行不是一回事。 single_step_run 不走引擎、不挂 Layer,所以没有额度检查、没有 OTel span、没有执行限制。它跑得通不代表全图跑得通。

  • 迭代/循环内部节点的输出不进草稿变量workflow_draft_variable_service.py:908-913),所以循环体里的节点做不到「引用上一次跑的上游值」这种单跑体验。

  • 草稿变量按 (app_id, user_id, node_id, name) 唯一。 同一个 app 里不同用户的调试互不干扰,但同一个用户多次调试会互相覆盖——最后一次跑赢。

  • 模型额度不足现在在单次模型调用处抛错reserve_model_quota_for_modelQuotaExceededErrorapi/core/app/llm/quota.py:139),这一次调用失败;旧版 LLMQuotaLayer「拿不到模型身份就中止整张图」的行为已随该 Layer 一起移除。

  • workflow_node_execution_offload 依赖 PostgreSQL「NULL 值互不相等」的默认唯一约束语义models/workflow.py:1223-1235)。换数据库要重新验证这个假设。


10. 代码地图(导航索引)

主题文件路径(克隆根相对)符号名
落库 Layerapi/core/app/workflow/layers/persistence.pyWorkflowPersistenceLayerPersistenceWorkflowInfo_update_node_execution_enqueue_trace_task
OTel 埋点 Layerapi/core/app/workflow/layers/observability.pyObservabilityLayer_NodeSpanContext
模型额度结算api/core/model_manager.pyQuotaManagedModelInstancereserve_quotarelease_quota_safely(配额函数在 api/core/app/llm/quota.py
限时 Layerapi/core/app/layers/timeslice_layer.pyTimeSliceLayer_checker_job
会话变量落库 Layerapi/core/app/layers/conversation_variable_persist_layer.pyConversationVariablePersistenceLayer
触发器日志 Layerapi/core/app/layers/trigger_post_layer.pyTriggerPostLayer_STATUS_MAP
Layer 挂载点(workflow)api/core/app/apps/workflow/app_runner.pyWorkflowAppRunner.runRedisChannel:164
Layer 挂载点(advanced chat)api/core/app/apps/advanced_chat/app_runner.pyAdvancedChatAppRunnerRedisChannel:215
引擎入口 / 内存通道api/core/workflow/workflow_entry.pyWorkflowEntry.__init__single_step_run
停止命令发送api/controllers/console/app/workflow.pyGraphEngineManager(...).send_stop_command
停止标志位(旧机制,01 章详述)api/core/app/apps/base_app_queue_manager.pyset_stop_flag_is_stoppedgenerate_task_stopped:{task_id}
运行/执行数据模型api/models/workflow.pyWorkflowRunWorkflowNodeExecutionModelWorkflowNodeExecutionOffloadWorkflowAppLogWorkflowArchiveLog
草稿变量模型api/models/workflow.pyWorkflowDraftVariableWorkflowDraftVariableFilenew_node_variable
写入侧仓储 + 卸载api/core/repositories/sqlalchemy_workflow_node_execution_repository.pysavesave_execution_data_truncate_and_upload
仓储工厂(写入侧)api/core/repositories/factory.pyDifyCoreRepositoryFactoryWorkflowNodeExecutionRepository
仓储工厂(读取侧)api/repositories/factory.pyDifyAPIRepositoryFactory
读取侧仓储接口api/repositories/api_workflow_run_repository.pyAPIWorkflowRunRepository
附加内容仓储api/repositories/execution_extra_content_repository.pyExecutionExtraContentRepository
草稿变量保存端口api/core/app/apps/draft_variable_saver.pyDraftVariableSaverDraftVariableSaverFactoryNoopDraftVariableSaver
草稿变量保存实现api/services/workflow_draft_variable_service.pyDraftVariableSaverDraftVarLoader_EXCLUDE_VARIABLE_NAMES_MAPPING_batch_upsert_draft_variable
保存器工厂分派api/core/app/apps/base_app_generator.py_get_draft_var_saver_factory_DebuggerDraftVariableSaver
单节点调试 APIapi/controllers/console/app/workflow.pyDraftWorkflowNodeRunApi
单节点调试服务api/services/workflow_service.pyrun_draft_workflow_node
单次迭代/循环入口api/core/app/apps/workflow/app_generator.pysingle_iteration_generatesingle_loop_generate
子图裁剪api/core/app/apps/workflow_app_runner.py_prepare_single_node_execution_get_graph_and_variable_pool_for_single_node_run
第三方 trace 队列api/core/ops/ops_trace_manager.pyTraceQueueManagerTraceTasksend_to_celery
前端 inspector 频道api/services/workflow/inspector_events.pypublish_node_changedpublish_workflow_completed
仓储/截断配置api/configs/feature/__init__.pyCORE_WORKFLOW_NODE_EXECUTION_REPOSITORYWORKFLOW_VARIABLE_TRUNCATION_MAX_SIZE
落库 Layer 单测(钩子签名)api/tests/unit_tests/core/app/workflow/test_persistence_layer.pylayer.initialize(read_only_state, command_channel=None)

相关章节

  • 01 一次运行的生命周期 —— 本章的事件流下游:SSE 管线怎么把 GraphEngineEvent 变成前端能看的流;其 §7.2 讲「停止」的队列侧旧机制,与本章 §4 的命令通道互为两半。
  • 03 从 JSON 到可执行图 —— 本章挂 Layer 的那个 GraphEngine 和它的图是怎么造出来的。
  • 05 停下来等人 —— 暂停/恢复相关的 PauseStatePersistenceLayerSuspendLayerworkflow_pauses 表,本章不覆盖。
  • 06 触发器与插件运行时 —— TriggerPostLayer 回写的 workflow_trigger_logs 从哪来。