跳到主要内容

数据截至 (上游 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
WorkflowEntryGraph 和一堆 Layer 拼成一台 GraphEngineapi/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

主线走一遍(高层,不进代码):

  1. 接住请求。 /v1/workflows/run 收到 JSON,按 app 模式分发给对应的 Generator,路上先过配额和限流。
  2. 一分为二。 Generator 建好队列管理器,另起一个线程去跑引擎,自己留在 HTTP 线程上等着从队列里拿事件。
  3. 备料。 后台线程从库里查出 Workflow,把系统变量、环境变量、用户输入灌进一个变量池,再把 graph JSON 交给节点工厂装配成图。
  4. 挂钩子。 落库、计费、可观测、暂停持久化这些横切能力,全部以 Layer 的形式挂到引擎上——引擎本身不认识数据库。
  5. 跑。 引擎逐节点执行,每产生一个事件,Layer 先看一遍,再被转成队列消息发往 HTTP 线程。
  6. 吐。 HTTP 线程那边的 TaskPipeline 把队列事件翻译成对外的响应对象;streaming 就逐条 data: {...} 推走,blocking 就攒到终态再一次性返回。
  7. (可选)停。 如果图里有人工输入节点,引擎会发出"暂停"事件:运行态被序列化成一行字符串存库,请求先返回;等人提交表单后,一个 Celery 任务把状态反序列化回来,走同一条 _generate 路径续跑。

3. 阅读地图(建议顺序)

Dify 的后端很大,本组文档挑出**"一张图如何被跑起来"**这条主线拆成 6 章,由浅入深:

  1. 一次运行的生命周期:从 HTTP 请求到 SSE 流(先读)。请求怎么进来、为什么要开两个线程、queue.Queue 两端各是谁、blocking 与 streaming 为什么只差一个布尔值、超时和"停止"信号从哪来。读完你能在脑子里画出一次运行的调用栈。
  2. 编辑侧:画布画出什么、数据库存什么、DSL 带走什么。前端把 nodes/edges/viewport 组装成什么样的 payload、_ 前缀的临时字段怎么被剥掉、hash 乐观锁怎么防并发覆盖、草稿与已发布版本的关系、DSL 导出时怎么顺带算出插件依赖。
  3. 从 JSON 到可执行图:graphon 边界与 DifyNodeFactory(核心)。哪些东西已经搬去了 graphon、Dify 手里还剩什么、节点注册表如何把内置节点和工作流本地节点合成一张表、版本不匹配时如何回退到最新实现、根节点怎么被推断出来。
  4. 执行期横切:Layer 钩子、留痕落库与单节点调试。Layer 是 Dify 唯一的横切扩展点:落库、LLM 配额、可观测、时间片、触发器回写各自监听哪些事件;以及"只跑一个节点"的调试通道如何绕开整张图。停止命令走的那条 Redis 命令通道也在这一章(第 1 章讲的是它的另一半——Redis 停止标志位)。
  5. 停下来等人:暂停、审批表单与恢复执行。暂停的本质是把 GraphRuntimeState 序列化成字符串;表单 token 按"接收方类型"分级发放,控制台、Web App、Service API 能看到的东西不一样;恢复走的是 Celery + 同一条生成路径。
  6. 入口的另一半与外部能力:触发器与插件运行时。除了人点"运行",还有 Webhook、定时、插件事件三种入口,它们统一落到异步队列里执行;以及节点里的工具/模型/触发器能力如何通过 HTTP 打到插件守护进程。

想最快抓住精华:读第 1 章的"两个线程 + 一个队列"和第 3 章的"graphon 边界",就掌握了这套架构 70% 的形状。

4. 巧妙之处(可借鉴的技术)

① 把图调度器整个搬出仓库,只留三道接缝。 GraphGraphEngineVariablePoolGraphRuntimeState、Layer 基类、内置节点,全部来自第三方包 graphon==0.7.0api/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 不一致就抛 WorkflowHashNotEqualErrorapi/services/workflow_service.py:437-438),前端据此刷新,避免两个标签页互相覆盖。

⑥ 节点版本"匹配不上就退到最新"。 matched_node_class or latest_node_classapi/core/workflow/node_factory.py:161)——老 DSL 里写着的节点版本如果已经没有对应实现,不是报错,而是回落到该类型的最新实现。代价是行为可能悄悄变化,收益是几年前导出的 DSL 今天还打得开

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

按"一次运行的顺序"排列。锚点列的行号 as-of 49a92f0;行号会漂移,优先用符号名 grep

主题文件路径符号名锚点
对外运行接口(/v1)api/controllers/service_api/app/workflow.pyWorkflowRunApi.postapi/controllers/service_api/app/workflow.py:326
按 app 模式分发 + 限流配额api/services/app_generate_service.pyAppGenerateService.generateapi/services/app_generate_service.py:88
起线程、建队列、返回流api/core/app/apps/workflow/app_generator.pyWorkflowAppGenerator._generateapi/core/app/apps/workflow/app_generator.py:321
后台线程入口api/core/app/apps/workflow/app_generator.py_generate_workerapi/core/app/apps/workflow/app_generator.py:616
恢复一次已暂停的运行api/core/app/apps/workflow/app_generator.pyWorkflowAppGenerator.resumeapi/core/app/apps/workflow/app_generator.py:274
备料并驱动引擎api/core/app/apps/workflow/app_runner.pyWorkflowAppRunner.runapi/core/app/apps/workflow/app_runner.py:74
引擎事件 → 队列消息api/core/app/apps/workflow_app_runner.pyWorkflowBasedAppRunner._handle_eventapi/core/app/apps/workflow_app_runner.py:409
队列消费循环(含超时/心跳)api/core/app/apps/base_app_queue_manager.pyAppQueueManager.listenapi/core/app/apps/base_app_queue_manager.py:64
队列事件 → 对外响应api/core/app/apps/workflow/generate_task_pipeline.pyWorkflowAppGenerateTaskPipeline.processapi/core/app/apps/workflow/generate_task_pipeline.py:126
SSE 行格式 data: {...}api/core/app/apps/base_app_generator.pyconvert_to_event_streamapi/core/app/apps/base_app_generator.py:313
Flask 侧 event-stream 响应api/libs/helper.pycompact_generate_responseapi/libs/helper.py:412
组装 GraphEngine + 挂内置 Layerapi/core/workflow/workflow_entry.pyWorkflowEntry.__init__api/core/workflow/workflow_entry.py:92-175
跑图并过响应流过滤器api/core/workflow/workflow_entry.pyWorkflowEntry.runiter_dify_graph_engine_eventsapi/core/workflow/workflow_entry.py:177:49
单节点调试执行api/core/workflow/workflow_entry.pysingle_step_runrun_free_nodeapi/core/workflow/workflow_entry.py:192:408
节点注册表(内置 + 本地合流)api/core/workflow/node_factory.pyregister_nodesget_node_type_classes_mappingapi/core/workflow/node_factory.py:120:117
版本解析与回退api/core/workflow/node_factory.pyresolve_workflow_node_classapi/core/workflow/node_factory.py:138
推断根节点api/core/workflow/node_factory.pyget_default_root_node_idapi/core/workflow/node_factory.py:172
JSON → Node 实例api/core/workflow/node_factory.pyDifyNodeFactory.create_nodeapi/core/workflow/node_factory.py:405
工作流本地节点(非 graphon 内置)api/core/workflow/nodes/agentagent_v2datasourceknowledge_indexknowledge_retrievaltrigger_plugintrigger_scheduletrigger_webhook目录,8 个子包
引擎外包声明api/pyproject.tomlgraphon==0.7.0api/pyproject.toml:48
运行/节点执行落库api/core/app/workflow/layers/persistence.pyWorkflowPersistenceLayer.on_eventapi/core/app/workflow/layers/persistence.py:83-133
暂停态序列化落库api/core/app/layers/pause_state_persist_layer.pyPauseStatePersistenceLayerWorkflowResumptionContextapi/core/app/layers/pause_state_persist_layer.py:77:36
异步执行的时间片与回写 Layerapi/core/app/layers/TimeSliceLayerTriggerPostLayerapi/core/app/layers/timeslice_layer.py:16trigger_post_layer.py:24
会话变量落库(chatflow 专用)api/core/app/layers/conversation_variable_persist_layer.pyConversationVariablePersistenceLayer挂载点 api/core/app/apps/advanced_chat/app_runner.py:283
图与特性的存储列api/models/workflow.pyWorkflow.graphWorkflow.graph_dictWorkflow.unique_hashapi/models/workflow.py:239:296:530
暂停记录表api/models/workflow.pyWorkflowPauseWorkflowPauseReasonapi/models/workflow.py:2117:2100
草稿变量表(调试用)api/models/workflow.pyWorkflowDraftVariableapi/models/workflow.py:1559
前端草稿 payload 组装web/app/components/workflow-app/hooks/use-nodes-sync-draft.tsgetPostParamsweb/app/components/workflow-app/hooks/use-nodes-sync-draft.ts:44-122
草稿落库 + 乐观锁api/services/workflow_service.pyWorkflowService.sync_draft_workflowapi/services/workflow_service.py:400
图结构校验api/services/workflow_service.pyvalidate_graph_structureapi/services/workflow_service.py:1796
发布为正式版本api/services/workflow_service.pyWorkflowService.publish_workflowapi/services/workflow_service.py:678
单节点调试入口api/services/workflow_service.pyrun_draft_workflow_nodeapi/services/workflow_service.py:1136
DSL 导出 / 导入 / 依赖抽取api/services/app_dsl_service.pyexport_dslimport_app_extract_dependencies_from_workflow_graphapi/services/app_dsl_service.py:664:89:645
人工输入表单提交与恢复排队api/services/human_input_service.pyHumanInputService.submit_form_by_tokenenqueue_resumeapi/services/human_input_service.py:195:265
表单可见性按调用面分级api/core/workflow/human_input_policy.pyHumanInputSurfacedisposition_for_surfaceapi/core/workflow/human_input_policy.py:19:74
恢复执行的 Celery 任务api/tasks/async_workflow_tasks.pyresume_workflow_executionapi/tasks/async_workflow_tasks.py:204
触发器统一入口(/triggers)api/controllers/trigger/trigger.pytrigger_endpointapi/controllers/trigger/trigger.py:17-18
触发事件分发api/services/trigger/trigger_service.pyTriggerService.process_endpointinvoke_trigger_eventapi/services/trigger/trigger_service.py:77:44
Webhook 请求解析与校验api/services/trigger/webhook_service.pyWebhookService.extract_and_validate_webhook_dataapi/services/trigger/webhook_service.py:174
异步执行队列(按套餐分队)api/tasks/async_workflow_tasks.pyexecute_workflow_professional / _team / _sandboxapi/tasks/async_workflow_tasks.py:54:70:86
插件守护进程客户端api/core/plugin/impl/base.pyplugin_daemon_inner_api_baseurl_request_with_plugin_daemon_response_streamapi/core/plugin/impl/base.py:49:299
触发器插件管理api/core/trigger/trigger_manager.pyTriggerManager.invoke_trigger_eventsubscribe_triggerapi/core/trigger/trigger_manager.py:150:198

两处诚实说明:

  • graphon 的源码不在本克隆内find 全仓无 graphon 目录),所以本组文档对引擎内部调度算法不做断言,只描述 Dify 一侧的调用与边界;凡涉及引擎行为的表述均以 Dify 的调用点为准。
  • SuspendLayerapi/core/app/layers/suspend_layer.py:7)在本 commit 的生产代码里没有挂载点,只被单元测试引用(api/tests/unit_tests/core/app/layers/test_suspend_layer.py)——列在这里是为了避免读者误以为它参与了暂停链路。