数据截至 (上游 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 就是这么让引擎暂停的。
一个必须先说清的边界: GraphEngine、GraphEngineLayer、RedisChannel、InMemoryChannel 都不在 Dify 仓库里,它们来自外部依赖 graphon==0.7.0(api/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_tokens和outputs的(api/core/app/layers/trigger_post_layer.py:75-79,on_event)。self.command_channel—— 反向命令通道,见 §4。
3.3 Dify 的 Layer 全家福
Dify 侧一共五个业务 Layer,分两个目录放,划分依据是「依赖多重」:
| Layer | 文件(api/ 下) | 干什么 | 挂在哪 |
|---|---|---|---|
WorkflowPersistenceLayer | core/app/workflow/layers/persistence.py:83 | 把运行/节点执行写进数据库,顺带投递 trace 任务 | app_runner 里手动挂 |
ObservabilityLayer | core/app/workflow/layers/observability.py:43 | 每个节点开一个 OpenTelemetry span | WorkflowEntry.__init__ 按开关挂 |
TimeSliceLayer | core/app/layers/timeslice_layer.py:16 | 定时检查是否超配额,超了就发 PAUSE | 异步触发任务传入 |
ConversationVariablePersistenceLayer | core/app/layers/conversation_variable_persist_layer.py:24 | 把 conversation.* 变量更新落库 | 只有 advanced-chat 挂 |
TriggerPostLayer | core/app/layers/trigger_post_layer.py:24 | 终态时回写 workflow_trigger_logs | 异步触发任务传入 |
(另有 PauseStatePersistenceLayer(core/app/layers/pause_state_persist_layer.py:77)与 SuspendLayer(core/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-284(TimeSliceLayer + 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_started(persistence.py:226):构造一个WorkflowNodeExecution领域对象(状态RUNNING),存进内存缓存,repository.save(...)落一条「跑起来了」的记录,再往前端的 inspector 频道推一条status="running"。_update_node_execution(persistence.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 就不挂ObservabilityLayer(api/core/workflow/workflow_entry.py:173-175的开关判断),都是删一行的事。 - 测试便宜。 单测直接
layer.initialize(fake_state, command_channel=None)然后喂事件,不用起引擎。
代价也是真的:
- 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)。 - 钩子不够用时只能捅私有方法——Dify 的解法是把这类职责搬出 Layer。 旧版
LLMQuotaLayer要在节点跑之前就让它失败,但 graphon 没给「跑前失败/跳过」的公开钩子,于是它直接把node._run换成一个返回FAILED的闭包并留 TODO 等上游开放公开钩子。上游后来把这一层整个移除,配额改为在模型调用边界结算(§4.3)——靠私有 override 顶住的需求,最终以「换挂载点」而不是「等钩子」收场。这是「扩展点边界画得不够宽」时最常见的症状与出路。 - 顺序耦合。 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 / 载体 | 适用 |
|---|---|---|---|
RedisChannel | api/core/app/apps/workflow/app_runner.py:164;api/core/app/apps/advanced_chat/app_runner.py:234 | f"workflow:{task_id}:commands" | 生产运行,Web 进程和执行进程可能不同 |
InMemoryChannel | api/core/workflow/workflow_entry.py:135(缺省) | 进程内对象 | 单节点调试、以及未显式传通道的一切入口 |
两个 app_runner 都以 Redis 通道为底,但写法已经分化:workflow 侧把 key 生成收进 app_task_command_channel_key(api/core/app/apps/workflow/app_runner.py:156-164),advanced-chat 侧则把 RedisChannel 和一个监听 Celery warm-shutdown 信号的 CelerySignalCommandChannel 组合成 CombinedCommandChannel(api/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:1201、api/controllers/web/workflow.py:131、api/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 设计里最巧的一环:它既是观察者,又能反向干预。
| 发送方 | 命令 | 触发条件 | 代码 |
|---|---|---|---|
| 停止 API | stop | 用户点击 | controllers/console/app/workflow.py:1201 |
TimeSliceLayer | CommandType.PAUSE | APScheduler 定时器发现资源配额到顶 | 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 已移除,配额改由 QuotaManagedModelInstance(api/core/model_manager.py:446)负责——invoke_llm 先 _reserve_quota_for_request 预留,拿到响应后按 usage 提交,finally 里 release_quota_safely 释放(model_manager.py:551-567);流式调用分「边收边提交」与「攒完再交」两种模式(_invoke_llm_stream,model_manager.py:569 起)。轮询式 LLM 的结算在 workflow 侧的 DifyPreparedPollingLLM._settle_polling_quota(api/core/workflow/node_runtime.py:332)。额度不足时 reserve_model_quota_for_model 直接抛 QuotaExceededError(api/core/app/llm/quota.py:139)——失败落在这一次模型调用上,而不是像旧 Layer 那样反向发 AbortCommand 中止整张图。
5. 留痕:跑完在数据库里剩下什么
5.1 五张表,各管一段
| 表(模型) | api/models/workflow.py | 一行代表 | 谁写的 |
|---|---|---|---|
workflow_runs(WorkflowRun) | :743 | 一次完整运行:状态、耗时、token、输入输出 | WorkflowPersistenceLayer 经仓储 |
workflow_node_executions(WorkflowNodeExecutionModel) | :922 | 一个节点的一次执行 | 同上 |
workflow_node_execution_offload(WorkflowNodeExecutionOffload) | :1171 | 某次执行被卸载到对象存储的大字段 | 仓储在 save_execution_data 里 |
workflow_app_logs(WorkflowAppLog) | :1280 | 面向「应用日志」列表的一条记录(不含调试) | SSE 管线,generate_task_pipeline.py:766 |
workflow_archive_logs(WorkflowArchiveLog) | :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_runs(triggered_from=DEBUGGING,见 api/models/enums.py:24-31),但不会污染应用日志列表。
WorkflowArchiveLog 则是反范式的快照:它把 WorkflowRun 的字段加 run_ 前缀(run_status / run_elapsed_time / run_total_tokens…)、把 WorkflowAppLog 的字段加 log_ 前缀,再塞一个 trigger_metadata(models/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_tokens、runtime_state.node_run_steps、runtime_state.exceptions_count。
5.3 两段式写入:save 和 save_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_upload(sqlalchemy_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 取附加内容)
| 层 | 工厂 | 接口形态 | 典型方法 |
|---|---|---|---|
| 写入侧 | DifyCoreRepositoryFactory(core/repositories/factory.py:55) | 极窄:save / save_execution_data / get_by_workflow_execution | 引擎线程里调 |
| 读取侧 | DifyAPIRepositoryFactory(repositories/factory.py:17,继承前者) | 很宽:分页、按时间批量取、软删、归档、暂停记录 | get_paginated_workflow_runs、create_archive_logs、get_expired_runs_batch |
实现类是配置字符串,不是硬编码。 工厂用 import_string(class_path) 动态加载(core/repositories/factory.py:86-94),路径来自四个配置项(api/configs/feature/__init__.py:955-980):
| 配置项 | 默认实现 | 换成什么 |
|---|---|---|
CORE_WORKFLOW_EXECUTION_REPOSITORY | SQLAlchemyWorkflowExecutionRepository | CeleryWorkflowExecutionRepository |
CORE_WORKFLOW_NODE_EXECUTION_REPOSITORY | SQLAlchemyWorkflowNodeExecutionRepository | CeleryWorkflowNodeExecutionRepository |
API_WORKFLOW_NODE_EXECUTION_REPOSITORY | DifyAPISQLAlchemyWorkflowNodeExecutionRepository | 自定义 |
API_WORKFLOW_RUN_REPOSITORY | DifyAPISQLAlchemyWorkflowRunRepository | 自定义 |
Celery 版的意义:把落库这个阻塞操作从引擎线程挪到后台 worker,代价是「刚写的立刻读」需要一层内存缓存兜着(api/core/repositories/celery_workflow_node_execution_repository.py:38-44 的类注释)。这也是 §3.5 那条「Layer 阻塞引擎线程」的官方解药。
ExecutionExtraContentRepository(repositories/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_variables(WorkflowDraftVariable) | models/workflow.py:1559 | 一个变量的当前草稿值 |
workflow_draft_variable_files(WorkflowDraftVariableFile) | models/workflow.py:2029 | 被卸载的大变量的元数据(size / length / 原始 value_type) |
WorkflowDraftVariable 上有几个非常「产品化」的字段:
node_id是个复用字段:普通节点填节点 id,会话变量填conversation,系统变量填sys(models/workflow.py:1624-1629的注释)。visible决定要不要在变量检查面板里显示,editable决定用户能不能改(:1573-1579)。IF_ELSE节点的输出一律不可见、不可编辑的系统变量一律不可见(services/workflow_draft_variable_service.py:1191-1196,_should_variable_be_visible)。last_edited_at为None表示「创建后没被人改过」(:1529-1535)。file_id非空表示这个值被卸载了,value里是截断版(:1567-1571)。
构造必须走三个工厂方法 new_conversation_variable / new_sys_variable / new_node_variable(models/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 |
NoopDraftVariableSaver | core/app/apps/draft_variable_saver.py:31 | 什么都不做 |
_DebuggerDraftVariableSaver | core/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_ASSIGNER | process_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
DraftVarLoader(services/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_from是SINGLE_STEP(models/workflow.py:965)。- 跑完还要再存一次草稿变量(
services/workflow_service.py:1246-1256),这样下游节点下次单跑时就能引用到本次的输出。闭环就是这么合上的。
6.6 单次迭代 / 单次循环:这个走引擎
「单独跑一次迭代」和单节点不一样——迭代体里可能有好几个节点,必须真的调度。所以它走的是完整的 generator 路径,只是把图裁小了。
入口:WorkflowAppGenerator.single_iteration_generate(core/app/apps/workflow/app_generator.py:431)和 single_loop_generate(:495),advanced-chat 有对应的一对(core/app/apps/advanced_chat/app_generator.py:320 / :391)。它们做的事:
- 造一个
invoke_from=DEBUGGER的WorkflowAppGenerateEntity,带上SingleIterationRunEntity(node_id, inputs)(app_generator.py:447-449)。 - 仓储用
WorkflowRunTriggeredFrom.DEBUGGING+WorkflowNodeExecutionTriggeredFrom.SINGLE_STEP(:461-472)。 - 造
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_run) | 否 | 否 | NULL | 只有那个 Node 对象 |
| 单次迭代 / 单次循环 | 是 | 是 | 有 | 裁剪后的子图 |
| 完整调试运行 | 是 | 是 | 有 | 全图 |
7. 可观测:两条互不相干的 trace 通路
新手最容易搞混的地方:Dify 里「trace」有两套,目标不同、路径不同、开关不同。
ObservabilityLayer | TraceQueueManager | |
|---|---|---|
| 面向谁 | 运维(APM) | 应用开发者(LLM 可观测平台) |
| 协议 | OpenTelemetry | Langfuse / 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_start(core/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_manager(persistence.py:92),图跑到终态时调 _enqueue_trace_task(persistence.py:476)造一个 TraceTask(TraceTaskName.WORKFLOW_TRACE, ...) 丢进去。
TraceQueueManager(core/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-1509;send_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. 巧妙之处(可以直接借鉴的)
-
生命周期钩子只管内存,事件才是事实来源。
WorkflowPersistenceLayer.on_graph_start只清缓存、on_graph_end直接return(persistence.py:111、:143)。所有落库都由事件驱动,于是「暂停后恢复」这种半程场景不需要为钩子写特例。 -
空对象消灭调用点的 if。
NoopDraftVariableSaver(core/app/apps/draft_variable_saver.py:31)让 SSE 管线无脑调saver.save(...),是否调试模式的判断被收敛到工厂一处(base_app_generator.py:332)。 -
可观测性优先于 I/O 效率的显式取舍。 inputs/outputs 本可以合并成一次卸载,Dify 选择分开,理由是合并会让节点在完成前不可观测——注释写了 20 行来解释(
models/workflow.py:1238-1259)。把「为什么没做那个显然的优化」写下来,比优化本身更有价值。 -
Layer 既是观察者又是干预者。
TimeSliceLayer通过self.command_channel反向发 PAUSE(timeslice_layer.py:54)。扩展点给了双向能力,避免了「为了停机再开一个后门」。 -
channel key 只由 task_id 决定。
f"workflow:{task_id}:commands"(app_runner.py:149)—— 停止请求不需要知道运行在哪台机器上,天然支持水平扩容。 -
仓储实现是配置字符串。 同一份 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_model抛QuotaExceededError(api/core/app/llm/quota.py:139),这一次调用失败;旧版LLMQuotaLayer「拿不到模型身份就中止整张图」的行为已随该 Layer 一起移除。 -
workflow_node_execution_offload依赖 PostgreSQL「NULL 值互不相等」的默认唯一约束语义(models/workflow.py:1223-1235)。换数据库要重新验证这个假设。
10. 代码地图(导航索引)
| 主题 | 文件路径(克隆根相对) | 符号名 |
|---|---|---|
| 落库 Layer | api/core/app/workflow/layers/persistence.py | WorkflowPersistenceLayer、PersistenceWorkflowInfo、_update_node_execution、_enqueue_trace_task |
| OTel 埋点 Layer | api/core/app/workflow/layers/observability.py | ObservabilityLayer、_NodeSpanContext |
| 模型额度结算 | api/core/model_manager.py | QuotaManagedModelInstance、reserve_quota、release_quota_safely(配额函数在 api/core/app/llm/quota.py) |
| 限时 Layer | api/core/app/layers/timeslice_layer.py | TimeSliceLayer、_checker_job |
| 会话变量落库 Layer | api/core/app/layers/conversation_variable_persist_layer.py | ConversationVariablePersistenceLayer |
| 触发器日志 Layer | api/core/app/layers/trigger_post_layer.py | TriggerPostLayer、_STATUS_MAP |
| Layer 挂载点(workflow) | api/core/app/apps/workflow/app_runner.py | WorkflowAppRunner.run(RedisChannel 在 :164) |
| Layer 挂载点(advanced chat) | api/core/app/apps/advanced_chat/app_runner.py | AdvancedChatAppRunner(RedisChannel 在 :215) |
| 引擎入口 / 内存通道 | api/core/workflow/workflow_entry.py | WorkflowEntry.__init__、single_step_run |
| 停止命令发送 | api/controllers/console/app/workflow.py | GraphEngineManager(...).send_stop_command |
| 停止标志位(旧机制,01 章详述) | api/core/app/apps/base_app_queue_manager.py | set_stop_flag、_is_stopped、generate_task_stopped:{task_id} |
| 运行/执行数据模型 | api/models/workflow.py | WorkflowRun、WorkflowNodeExecutionModel、WorkflowNodeExecutionOffload、WorkflowAppLog、WorkflowArchiveLog |
| 草稿变量模型 | api/models/workflow.py | WorkflowDraftVariable、WorkflowDraftVariableFile、new_node_variable |
| 写入侧仓储 + 卸载 | api/core/repositories/sqlalchemy_workflow_node_execution_repository.py | save、save_execution_data、_truncate_and_upload |
| 仓储工厂(写入侧) | api/core/repositories/factory.py | DifyCoreRepositoryFactory、WorkflowNodeExecutionRepository |
| 仓储工厂(读取侧) | api/repositories/factory.py | DifyAPIRepositoryFactory |
| 读取侧仓储接口 | api/repositories/api_workflow_run_repository.py | APIWorkflowRunRepository |
| 附加内容仓储 | api/repositories/execution_extra_content_repository.py | ExecutionExtraContentRepository |
| 草稿变量保存端口 | api/core/app/apps/draft_variable_saver.py | DraftVariableSaver、DraftVariableSaverFactory、NoopDraftVariableSaver |
| 草稿变量保存实现 | api/services/workflow_draft_variable_service.py | DraftVariableSaver、DraftVarLoader、_EXCLUDE_VARIABLE_NAMES_MAPPING、_batch_upsert_draft_variable |
| 保存器工厂分派 | api/core/app/apps/base_app_generator.py | _get_draft_var_saver_factory、_DebuggerDraftVariableSaver |
| 单节点调试 API | api/controllers/console/app/workflow.py | DraftWorkflowNodeRunApi |
| 单节点调试服务 | api/services/workflow_service.py | run_draft_workflow_node |
| 单次迭代/循环入口 | api/core/app/apps/workflow/app_generator.py | single_iteration_generate、single_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.py | TraceQueueManager、TraceTask、send_to_celery |
| 前端 inspector 频道 | api/services/workflow/inspector_events.py | publish_node_changed、publish_workflow_completed |
| 仓储/截断配置 | api/configs/feature/__init__.py | CORE_WORKFLOW_NODE_EXECUTION_REPOSITORY、WORKFLOW_VARIABLE_TRUNCATION_MAX_SIZE |
| 落库 Layer 单测(钩子签名) | api/tests/unit_tests/core/app/workflow/test_persistence_layer.py | layer.initialize(read_only_state, command_channel=None) |
相关章节
- 01 一次运行的生命周期 —— 本章的事件流下游:SSE 管线怎么把
GraphEngineEvent变成前端能看的流;其 §7.2 讲「停止」的队列侧旧机制,与本章 §4 的命令通道互为两半。 - 03 从 JSON 到可执行图 —— 本章挂 Layer 的那个
GraphEngine和它的图是怎么造出来的。 - 05 停下来等人 —— 暂停/恢复相关的
PauseStatePersistenceLayer、SuspendLayer和workflow_pauses表,本章不覆盖。 - 06 触发器与插件运行时 ——
TriggerPostLayer回写的workflow_trigger_logs从哪来。