数据截至 (上游 commit b78a3462c9a6)
Dify — 可视化工作流平台:画布出图、引擎外包、执行可暂停
30 秒导读: Dify 让你在浏览器里拖一张流程图("取输入 → 调大模型 → 查知识库 → 输出"),然后把这张图变成一个可以被 HTTP 调用、边跑边吐字的线上服务。本组文档不讲怎么用 Dify,而是讲这张图在代码里到底经历了什么:它被存成一列 JSON、被一个工厂装配成节点对象、交给一个已经搬出仓库的外部引擎包
graphon去调度,运行中的每一个事件再穿过一串 Layer 钩子落库、计费、打点,最后变成data: {...}推给浏览器;中途还能停下来等人点"批准",几小时后从序列化的运行态原地续跑。
1. 这是什么(零基础也能懂)
一句话定义: Dify 是一个开源的 LLM 应用平台,其中最核心的一块是可视化工作流——你在画布上连出的那张图,会被平台当成一份可执行程序跑起来。
解决什么问题 / 给谁用:
假设你要做一个"客服工单助手":收到工单 → 判断是不是投诉 → 是就查知识库拟一份回复 → 拿给主管过目 → 主管点了批准才真发出去。用代码写,你要自己处理并发、超时、流式输出、中途暂停、错误重试、多租户隔离。Dify 的卖点是:这些你都不写,你只画图。
它把什么变成了产品能力:
| 你在画布上做的事 | 平台在背后给你的东西 |
|---|---|
| 拖节点、连线 | 一份可版本化的 graph JSON,存在 workflows 表的一列里 |
| 点"运行" | 一次多租户的后台执行 + 实时 SSE 事件流 |
| 放一个"人工输入"节点 | 运行会真的停下来,运行态被序列化落库,等表单提交后续跑 |
| 放一个"Webhook 触发器"节点 | 一个对外的 HTTP 端点,被打就异步排队跑这条流 |
| 点"导出 DSL" | 一份带依赖清单的 YAML,可以在另一个 Dify 实例里导入 |
用起来什么样(对外 API):
# 真实路由: /v1 前缀 + /workflows/run(api/controllers/service_api/__init__.py:6)
curl -X POST 'https://api.dify.ai/v1/workflows/run' \
-H 'Authorization: Bearer app-xxxxxx' \
-H 'Content-Type: application/json' \
-d '{
"inputs": {"query": "订单一直没发货"},
"response_mode": "streaming",
"user": "user-42"
}'
# streaming 模式下响应是 text/event-stream,一行行的 data: {...}
一句话直觉/类比: 把画布当源码编辑器,把 graph JSON 当编译产物,把 graphon 当虚拟机——Dify 本体则是那个把源码取出来、装配成字节码、喂给虚拟机、再把虚拟机吐出的每一条日志转发给用户和数据库的运行时宿主。
2. 顶层全景(它大概怎么转)
整个系统被一条缝清楚地切成两半:编辑期(慢、事务性、和人打交道)和运行期(快、流式、和模型打交道)。两边唯一的交接物就是一张 graph JSON。
怎么读这张图: 上半是编辑期、下半是运行期,中间那条横线是两者的唯一交接口——workflows.graph 这一列文本。箭头方向就是数据流向。
── 编辑期(慢、事务性)─────────────────────────────────
┌──────────────┐ POST 草稿(带 hash) ┌──────────────┐
│ 画布 (React) │ ────────────────────→ │ DSL (YAML) │
│ nodes / edges │ ←── 导入 ── 导出 ──→ │ 带依赖清单 │
└──────┬───────┘ └──────────────┘
│ 落库
▼
╔══════════════════════════════╗
║ workflows 表 · graph 列(JSON) ║ ← 编辑期与运行期的唯一交接物
╚══════════════╤═══════════════╝
│ graph_dict
── 运行期(快、流式)─┼──────────────────────────────────
▼
┌──────────────────┐ ┌─────────────────────┐
│ DifyNodeFactory │ ───→ │ graphon.GraphEngine │ ← 外部包,只管调度
│ JSON → Node 实例 │ └──────────┬──────────┘
└──────────────────┘ │ 事件流
▼
┌────────────────────────────┐
│ Layer 钩子 → 队列 → Pipeline │ → SSE / 数据库
└────────────────────────────┘
部件一句话职责:
| 部件 | 干什么 | 在哪个文件(克隆根相对) |
|---|---|---|
WorkflowAppGenerator | 一次运行的总装配工:建队列、起线程、返回流 | api/core/app/apps/workflow/app_generator.py |
WorkflowAppRunner | 在后台线程里备料(变量池、图)并驱动引擎 | api/core/app/apps/workflow/app_runner.py |
WorkflowEntry | 把 Graph 和一堆 Layer 拼成一台 GraphEngine | api/core/workflow/workflow_entry.py |
DifyNodeFactory | 把一个节点的 JSON 配置实例化成 Node 对象 | api/core/workflow/node_factory.py |
graphon(外部包) | 图、引擎、变量池、内置节点、Layer 基类 | 第三方依赖,graphon==0.7.0 |
AppQueueManager | 引擎线程与 HTTP 线程之间的那根管子 | api/core/app/apps/base_app_queue_manager.py |
WorkflowAppGenerateTaskPipeline | 把队列事件翻译成对外的流式/阻塞响应 | api/core/app/apps/workflow/generate_task_pipeline.py |
WorkflowPersistenceLayer | 监听引擎事件,把运行与节点执行写进库 | api/core/app/workflow/layers/persistence.py |
PauseStatePersistenceLayer | 收到"暂停"事件就把运行态序列化存库 | api/core/app/layers/pause_state_persist_layer.py |
Workflow 模型 | graph / features / 环境变量都是文本列 | api/models/workflow.py |
AppDslService | 导出/导入 YAML,并抽出插件依赖清单 | api/services/app_dsl_service.py |
主线走一遍(高层,不进代码):
- 接住请求。
/v1/workflows/run收到 JSON,按 app 模式分发给对应的 Generator,路上先过配额和限流。 - 一分为二。 Generator 建好队列管理器,另起一个线程去跑引擎,自己留在 HTTP 线程上等着从队列里拿事件。
- 备料。 后台线程从库里查出
Workflow,把系统变量、环境变量、用户输入灌进一个变量池,再把graphJSON 交给节 点工厂装配成图。 - 挂钩子。 落库、计费、可观测、暂停持久化这些横切能力,全部以 Layer 的形式挂到引擎上——引擎本身不认识数据库。
- 跑。 引擎逐节点执行,每产生一个事件,Layer 先看一遍,再被转成队列消息发往 HTTP 线程。
- 吐。 HTTP 线程那边的 TaskPipeline 把队列事件翻译成对外的响应对象;streaming 就逐条
data: {...}推走,blocking 就攒到终态再一次性返回。 - (可选)停。 如果图里有人工输入节点,引擎会发出"暂停"事件:运行态被序列化成一行字符串存库,请求先返回;等人提交表单后,一个 Celery 任务把状态反序列化回来,走同一条
_generate路径续跑。
3. 阅读地图(建议顺序)
Dify 的后端很大,本组文档挑出**"一张图如何被跑起来"**这条主线拆成 6 章,由浅入深:
- 一次运行的生命周期:从 HTTP 请求到 SSE 流(先读)。请求怎么进来、为什么要开两个线程、
queue.Queue两端各是谁、blocking 与 streaming 为什么只差一个布尔值、超时和"停止"信号从哪来。读完你能在脑子里画出一次运行的调用栈。 - 编辑侧:画布画出什么、数据库存什么、DSL 带走什么。前端把
nodes/edges/viewport组装成什么样的 payload、_前缀的临时字段怎么被剥掉、hash乐观锁怎么防并发覆盖、草稿与已发布版本的关系、DSL 导出时怎么顺带算出插件依赖。 - 从 JSON 到可执行图:graphon 边界与 DifyNodeFactory(核心)。哪些东西已经搬去了
graphon、Dify 手里还剩什么、节点注册表如何把内置节点和工作流本地节点合成一张表、版本不匹配时如何回退到最新实现、根节点怎么被推断出来。 - 执行期横切:Layer 钩子、留痕落库与单节点调试。Layer 是 Dify 唯一的横切扩展点:落库、LLM 配额、可观测、时间片、触发器回写各自监听哪些事件;以及"只跑一个节点"的调试通道如何绕开整张图。停止命令走的那条 Redis 命令通道也在这一章(第 1 章讲的是它的另一半——Redis 停止标志位)。
- 停下来等人:暂停、审批表单与恢复执行。暂停的本质是把
GraphRuntimeState序列化成字符串;表单 token 按"接收方类型"分级发放,控制台、Web App、Service API 能看到的东西不一样;恢复走的是 Celery + 同一条生成路径。 - 入口的另一半与外部能力:触发器与插件运行时。除了人点"运行",还有 Webhook、定时、插件事件三种入口,它们统一落到异步队列里执行;以及节点里的工具/模型/触发器能力如何通过 HTTP 打到插件守护进程。
想最快抓住精华:读第 1 章的"两个线程 + 一个队列"和第 3 章的"graphon 边界",就掌握了这套架构 70% 的形状。
4. 巧妙之处(可借鉴的技术)
① 把图调度器整个搬出仓库,只留三道接缝。
Graph、GraphEngine、VariablePool、GraphRuntimeState、Layer 基类、内置节点,全部来自第三方包 graphon==0.7.0(api/pyproject.toml:48),仓库里没有它的源码。Dify 自己只守住三处接缝:节点工厂(怎么把 JSON 变成对象)、Layer(横切怎么插进去)、命令通道(外面怎么喊停)。妙在这个切法让"调度算法"和"多租户业务"彻底分家——api/core/workflow/workflow_entry.py:142 那一处 GraphEngine(...) 构造,就是两边的全部交界面。
② 一次运行 = 两个线程 + 一个 queue.Queue,流式与阻塞共用一条代码路径。
Generator 在起线程之前先 db.session.close() 释放连接,再用 contextvars.copy_context() 把请求上下文复制给工作线程(api/core/app/apps/workflow/app_generator.py:321-395)。主线程只负责 listen() 拉队列。于是 blocking 和 streaming 的差别缩到只剩 _handle_response 里的一个 stream 布尔——两种模式跑的是同一段引擎代码。
③ 队列的 listen() 顺手兼职三件事。
同一个 while True 循环里,除了取消息,还负责:超过 APP_MAX_EXECUTION_TIME 就自己发一个停止事件、 检测到外部停止标志也发停止事件、每 10 秒补一个 ping 心跳防止连接被中间层掐断(api/core/app/apps/base_app_queue_manager.py:64-98)。把"超时/取消/保活"塞进消费循环,省掉了三个定时器。
④ 暂停 = 把运行态 dumps() 成一行字符串。
PauseStatePersistenceLayer 只监听一种事件 GraphRunPausedEvent,收到就把 graph_runtime_state.dumps() 和生成实体一起打包成 WorkflowResumptionContext 写库(api/core/app/layers/pause_state_persist_layer.py:77-159)。恢复时反序列化回来,调 WorkflowAppGenerator.resume——而 resume 内部只是给 _generate 多传一个 graph_runtime_state 参数(api/core/app/apps/workflow/app_generator.py:274-319)。"续跑"没有独立的执行路径,这是它敢在生产里暂停几小时的底气。
⑤ 草稿同步的两道洁癖。
前端在发草稿前,把节点/边 data 里所有 _ 开头的键就地删掉(web/app/components/workflow-app/hooks/use-nodes-sync-draft.ts:73-91)——UI 态(选中、临时、拖拽中)永远进不了数据库。同时带上一个 hash,服务端一旦发现和当前草稿的 unique_hash 不一致就抛 WorkflowHashNotEqualError(api/services/workflow_service.py:437-438),前端据此刷新,避免两个标签页互相覆盖。
⑥ 节点版本"匹配不上就退到最新"。
matched_node_class or latest_node_class(api/core/workflow/node_factory.py:161)——老 DSL 里写着的节点版本如果已经没有对应实现,不是报错,而是回落到该类型的最新实现。代价是行为可能悄悄变化,收益是几年前导出的 DSL 今天还打得开。
5. 代码地图(导航索引)
按"一次运行的顺序"排列。锚点列的行号 as-of 49a92f0;行号会漂移,优先用符号名 grep。
| 主题 | 文件路径 | 符号名 | 锚点 |
|---|---|---|---|
| 对外运行接口(/v1) | api/controllers/service_api/app/workflow.py | WorkflowRunApi.post | api/controllers/service_api/app/workflow.py:326 |
| 按 app 模式分发 + 限流配额 | api/services/app_generate_service.py | AppGenerateService.generate | api/services/app_generate_service.py:88 |
| 起线程、建队列、返回流 | api/core/app/apps/workflow/app_generator.py | WorkflowAppGenerator._generate | api/core/app/apps/workflow/app_generator.py:321 |
| 后台线程入口 | api/core/app/apps/workflow/app_generator.py | _generate_worker | api/core/app/apps/workflow/app_generator.py:616 |
| 恢复一次已暂停的运行 | api/core/app/apps/workflow/app_generator.py | WorkflowAppGenerator.resume | api/core/app/apps/workflow/app_generator.py:274 |
| 备料并驱动引擎 | api/core/app/apps/workflow/app_runner.py | WorkflowAppRunner.run | api/core/app/apps/workflow/app_runner.py:74 |
| 引擎事件 → 队列消息 | api/core/app/apps/workflow_app_runner.py | WorkflowBasedAppRunner._handle_event | api/core/app/apps/workflow_app_runner.py:409 |
| 队列消费循环(含超时/心跳) | api/core/app/apps/base_app_queue_manager.py | AppQueueManager.listen | api/core/app/apps/base_app_queue_manager.py:64 |
| 队列事件 → 对外响应 | api/core/app/apps/workflow/generate_task_pipeline.py | WorkflowAppGenerateTaskPipeline.process | api/core/app/apps/workflow/generate_task_pipeline.py:126 |
SSE 行格式 data: {...} | api/core/app/apps/base_app_generator.py | convert_to_event_stream | api/core/app/apps/base_app_generator.py:313 |
| Flask 侧 event-stream 响应 | api/libs/helper.py | compact_generate_response | api/libs/helper.py:412 |
| 组装 GraphEngine + 挂内置 Layer | api/core/workflow/workflow_entry.py | WorkflowEntry.__init__ | api/core/workflow/workflow_entry.py:92-175 |
| 跑图并过响应流过滤器 | api/core/workflow/workflow_entry.py | WorkflowEntry.run、iter_dify_graph_engine_events | api/core/workflow/workflow_entry.py:177、:49 |
| 单节点调试执行 | api/core/workflow/workflow_entry.py | single_step_run、run_free_node | api/core/workflow/workflow_entry.py:192、:408 |
| 节点注册表(内置 + 本地 合流) | api/core/workflow/node_factory.py | register_nodes、get_node_type_classes_mapping | api/core/workflow/node_factory.py:120、:117 |
| 版本解析与回退 | api/core/workflow/node_factory.py | resolve_workflow_node_class | api/core/workflow/node_factory.py:138 |
| 推断根节点 | api/core/workflow/node_factory.py | get_default_root_node_id | api/core/workflow/node_factory.py:172 |
| JSON → Node 实例 | api/core/workflow/node_factory.py | DifyNodeFactory.create_node | api/core/workflow/node_factory.py:405 |
| 工作流本地节点(非 graphon 内置) | api/core/workflow/nodes/ | agent、agent_v2、datasource、knowledge_index、knowledge_retrieval、trigger_plugin、trigger_schedule、trigger_webhook | 目录,8 个子包 |
| 引擎外包声明 | api/pyproject.toml | graphon==0.7.0 | api/pyproject.toml:48 |
| 运行/节点执行落库 | api/core/app/workflow/layers/persistence.py | WorkflowPersistenceLayer.on_event | api/core/app/workflow/layers/persistence.py:83-133 |
| 暂停态序列化落库 | api/core/app/layers/pause_state_persist_layer.py | PauseStatePersistenceLayer、WorkflowResumptionContext | api/core/app/layers/pause_state_persist_layer.py:77、:36 |
| 异步执行的时间片与回写 Layer | api/core/app/layers/ | TimeSliceLayer、TriggerPostLayer | api/core/app/layers/timeslice_layer.py:16、trigger_post_layer.py:24 |
| 会话变量落库(chatflow 专用) | api/core/app/layers/conversation_variable_persist_layer.py | ConversationVariablePersistenceLayer | 挂载点 api/core/app/apps/advanced_chat/app_runner.py:283 |
| 图与特性的存储列 | api/models/workflow.py | Workflow.graph、Workflow.graph_dict、Workflow.unique_hash | api/models/workflow.py:239、:296、:530 |
| 暂停记录表 | api/models/workflow.py | WorkflowPause、WorkflowPauseReason | api/models/workflow.py:2117、:2100 |
| 草稿变量表(调试用) | api/models/workflow.py | WorkflowDraftVariable | api/models/workflow.py:1559 |
| 前端草稿 payload 组装 | web/app/components/workflow-app/hooks/use-nodes-sync-draft.ts | getPostParams | web/app/components/workflow-app/hooks/use-nodes-sync-draft.ts:44-122 |
| 草稿落库 + 乐观锁 | api/services/workflow_service.py | WorkflowService.sync_draft_workflow | api/services/workflow_service.py:400 |
| 图结构校验 | api/services/workflow_service.py | validate_graph_structure | api/services/workflow_service.py:1796 |
| 发布为正式版本 | api/services/workflow_service.py | WorkflowService.publish_workflow | api/services/workflow_service.py:678 |
| 单节点调试入口 | api/services/workflow_service.py | run_draft_workflow_node | api/services/workflow_service.py:1136 |
| DSL 导出 / 导入 / 依赖抽取 | api/services/app_dsl_service.py | export_dsl、import_app、_extract_dependencies_from_workflow_graph | api/services/app_dsl_service.py:664、:89、:645 |
| 人工输入表单提交与恢复排队 | api/services/human_input_service.py | HumanInputService.submit_form_by_token、enqueue_resume | api/services/human_input_service.py:195、:265 |
| 表单可见性按调用面分级 | api/core/workflow/human_input_policy.py | HumanInputSurface、disposition_for_surface | api/core/workflow/human_input_policy.py:19、:74 |
| 恢复执行的 Celery 任务 | api/tasks/async_workflow_tasks.py | resume_workflow_execution | api/tasks/async_workflow_tasks.py:204 |
| 触发器统一入口(/triggers) | api/controllers/trigger/trigger.py | trigger_endpoint | api/controllers/trigger/trigger.py:17-18 |
| 触发事件分发 | api/services/trigger/trigger_service.py | TriggerService.process_endpoint、invoke_trigger_event | api/services/trigger/trigger_service.py:77、:44 |
| Webhook 请求解析与校验 | api/services/trigger/webhook_service.py | WebhookService.extract_and_validate_webhook_data | api/services/trigger/webhook_service.py:174 |
| 异步执行队列(按套餐分队) | api/tasks/async_workflow_tasks.py | execute_workflow_professional / _team / _sandbox | api/tasks/async_workflow_tasks.py:54、:70、:86 |
| 插件守护进程客户端 | api/core/plugin/impl/base.py | plugin_daemon_inner_api_baseurl、_request_with_plugin_daemon_response_stream | api/core/plugin/impl/base.py:49、:299 |
| 触发器插件管理 | api/core/trigger/trigger_manager.py | TriggerManager.invoke_trigger_event、subscribe_trigger | api/core/trigger/trigger_manager.py:150、:198 |
两处诚实说明:
graphon的源码不在本克隆内(find全仓无graphon目录),所以本组文档对引擎内部调度算法不做断言,只描述 Dify 一侧的调用与边界;凡涉及引擎行为的表述均以 Dify 的调用点为准。SuspendLayer(api/core/app/layers/suspend_layer.py:7)在本 commit 的生产代码里没有挂载点,只被单元测试引用(api/tests/unit_tests/core/app/layers/test_suspend_layer.py)——列在这里是为了避免读者误以为它参与了暂停链路。