Day 11 / 共 20 天 · 阶段3 工作流引擎

变量系统与变量池:工作流里的数据是怎么"存与取"的

工作流是一张图,节点一个接一个跑。可是「开始节点」收到的用户问题,怎么被后面的「LLM 节点」用到?「LLM 节点」的输出,又怎么被「回复节点」引用?答案就是今天的主角——变量池(VariablePool)。弄清三件事:①变量怎么用一个"抽屉编号"(selector)唯一定位;②系统变量、环境变量、会话变量、节点输出这四类怎么用命名空间前缀分开放;③开机时怎么把初始变量装进池、节点跑完又怎么把输出写回池。这是理解整个工作流数据流的地基。

📍 你在 20 天里的位置(阶段3:工作流引擎)
S1 起步 S2 模型运行时 D11 变量系统 D12 事件与流式 S4 RAG 知识库 S5 工具/Agent S6 收官
💡 先用两个类比兜住今天 类比一:变量池像公司里的一块共享白板。每个节点跑完,把自己的产出写到白板的某个格子上;后面的节点要用,就照着"格子地址"去读。谁都不直接给谁递纸条,全靠这块公共白板中转——解耦。类比二:selector(选择器)像文件柜的"抽屉编号",形如 [node_id, 字段名]——先找到哪个抽屉(哪个节点),再找抽屉里哪个文件夹(哪个字段)。而 sys / env / conversation 这些前缀,就像抽屉上贴的分区标签:系统区、环境区、会话区,各放各的,不会混。
L01

痛点:一张图里,节点之间怎么传数据?

🤔 痛点你在 Dify 画布上拉了三个节点:开始 → LLM → 直接回复。LLM 节点的 Prompt 里写了 {{#sys.query#}}(用户的问题),回复节点里写了 {{#llm.text#}}(大模型的输出)。问题来了:运行时这些 {{#...#}} 到底怎么被替换成真实值?节点之间又不是函数调用,谁也不认识谁,凭什么 LLM 能拿到"开始节点"收进来的问题?如果每个节点都直接持有前一个节点的引用,那图一复杂就成了一团乱麻。
💡 本质:一块共享的"键值白板"Dify 的解法是设一个全局变量池 VariablePool:本质就是一个大字典。key 是 selector(一个由字符串组成的元组,如 ("sys", "query")),value 是包装好的变量值(Segment)。节点跑完 → 把输出 add(selector, value) 写进池;下一个节点要用 → 按 {{#...#}} 解析出 selector → get(selector) 读出来。节点彼此不认识,只认这块白板。今天我们不看池的底层实现(它在依赖库 graphon 里),而是看 Dify 在 core/workflow/ 下"怎么定义变量、怎么装池、怎么读池"这一层。
L02

本质:selector 是"抽屉编号"

一个变量的地址(selector)永远是"至少两段":第一段是命名空间/节点 id,第二段是字段名。比如系统变量 query 的地址就是 ("sys", "query")。看最基础的一个辅助函数 system_variable_selectorcore/workflow/system_variables.py:59):

# core/workflow/system_variables.py:55
def system_variable_name(key: str | SystemVariableKey) -> str:
    return key.value if isinstance(key, SystemVariableKey) else key

def system_variable_selector(key: str | SystemVariableKey) -> tuple[str, str]:
    return SYSTEM_VARIABLE_NODE_ID, system_variable_name(key)   # ★ ("sys", "query")
system_variable_name把 key 归一成纯字符串。传进来的可能是枚举(SystemVariableKey.QUERY),也可能已经是字符串("query"),统一取出 "query"
system_variable_selector★拼出那个"抽屉编号":(SYSTEM_VARIABLE_NODE_ID, "query"),也就是 ("sys", "query")。整个系统里凡是要读写系统变量 query,都用这个函数拿地址——地址集中生成,绝不散落硬编码
大白话你在画布里写的 {{#sys.query#}},中间那个点两边正好对应 selector 的两段:sys 是抽屉、query 是文件。运行时引擎把这个字符串拆成 ("sys","query"),拿去变量池一查,就换成了真实的用户问题。所谓"变量引用",本质就是把字符串地址翻译成字典的 key 去查表
L03

四个命名空间前缀:抽屉上的分区标签

那些"第一段"从哪来?Dify 把它们收拢在一个极短的文件里 variable_prefixes.pycore/workflow/variable_prefixes.py:1):

# core/workflow/variable_prefixes.py:1
SYSTEM_VARIABLE_NODE_ID = "sys"
ENVIRONMENT_VARIABLE_NODE_ID = "env"
CONVERSATION_VARIABLE_NODE_ID = "conversation"
RAG_PIPELINE_VARIABLE_NODE_ID = "rag"
前缀装什么画布里写法
sys系统给的:用户问题、文件、会话 id、用户 id…(运行时自动注入){{#sys.query#}}
env环境变量:整个 app 级别的常量(如某个 API 地址){{#env.xxx#}}
conversation会话变量:多轮对话里能被读写、跨轮次保留的状态{{#conversation.xxx#}}
ragRAG 流水线专用变量(知识库工作流场景)
某节点 id该节点的输出(如 LLM 节点的 text{{#1711429….text#}}
💡 为什么把这四个常量单独抽一个文件?因为它们是"魔法字符串"——一旦某处写成 "sys"、另一处写成 "system",变量就永远读不到、还极难排查。把它们定义成唯一的常量并集中在一个文件,别的模块 import 来用,就杜绝了拼写漂移。这跟你项目里"把所有魔法数抽成常量"是同一个工程直觉,只是 Dify 把它做到了模块级别。
L04

SystemVariableKey:系统变量到底有哪些

sys 这个抽屉里具体能放什么?由一个枚举钉死 SystemVariableKeycore/workflow/system_variables.py:22):

# core/workflow/system_variables.py:22
class SystemVariableKey(StrEnum):
    QUERY = "query"                       # 用户这轮的问题
    FILES = "files"                       # 用户上传的文件
    CONVERSATION_ID = "conversation_id"   # 会话 id
    USER_ID = "user_id"                   # 用户 id
    DIALOGUE_COUNT = "dialogue_count"     # 第几轮对话
    APP_ID = "app_id"
    WORKFLOW_ID = "workflow_id"
    WORKFLOW_EXECUTION_ID = "workflow_run_id"   # 注意:枚举名和值故意不同
    ...
    INVOKE_FROM = "invoke_from"           # 从哪触发(调试/API/前端)

那把这些系统变量真正"造出来"的是 build_system_variablescore/workflow/system_variables.py:81):

# core/workflow/system_variables.py:81
def build_system_variables(values=None, /, **kwargs) -> list[Variable]:
    normalized = _normalize_system_variable_values(values, **kwargs)   # ① 清洗
    return [
        segment_to_variable(                                           # ② 逐个包装成 Variable
            segment=build_segment(value),
            selector=system_variable_selector(key),                    #    地址 = ("sys", key)
            name=key,
        )
        for key, value in normalized.items()
    ]
_normalize_system_variable_values清洗一遍(system_variables.py:63):None 的值直接丢弃、把旧字段名 workflow_execution_id 兼容映射到新值、并且兜底 files 至少是空列表setdefault(FILES, []))——所以就算用户没传文件,{{#sys.files#}} 也不会炸。
build_segment(value)把裸值(字符串/列表/数字)包装成 Segment——变量池里存的是带类型信息的 Segment,不是裸值。
selector = ("sys", key)★每个系统变量的地址都用 L02 那个函数生成,前缀统一是 sys
注意 WORKFLOW_EXECUTION_ID = "workflow_run_id"——枚举名字故意不一样。这是历史包袱:代码里想用新名字 workflow_execution_id,但存量数据和画布引用还叫 workflow_run_id,于是保留旧值。_normalize 里那段 pop("workflow_execution_id") 就是在做新旧名兼容。真实项目里这种"改名不改值"随处可见,别被它绕晕。
L05

build_bootstrap_variables:开机时给白板装料

工作流一启动,得先把"初始变量"一次性装进池。这活由 build_bootstrap_variablescore/workflow/system_variables.py:109)干——它把四类变量各自盖上正确的命名空间前缀:

# core/workflow/system_variables.py:109
def build_bootstrap_variables(*, system_variables=(), environment_variables=(),
                              conversation_variables=(), rag_pipeline_variables=()):
    variables = [
        *(_with_selector(v, SYSTEM_VARIABLE_NODE_ID)       for v in system_variables),        # sys.*
        *(_with_selector(v, ENVIRONMENT_VARIABLE_NODE_ID)  for v in environment_variables),   # env.*
        *(_with_selector(v, CONVERSATION_VARIABLE_NODE_ID) for v in conversation_variables),  # conversation.*
    ]
    # rag 变量按"归属节点"分组打包
    rag_pipeline_variables_map = defaultdict(dict)
    for rag_var in rag_pipeline_variables:
        node_id = rag_var.variable.belong_to_node_id
        rag_pipeline_variables_map[node_id][rag_var.variable.variable] = rag_var.value
    for node_id, value in rag_pipeline_variables_map.items():
        variables.append(segment_to_variable(
            segment=build_segment(value),
            selector=(RAG_PIPELINE_VARIABLE_NODE_ID, node_id), name=node_id))
    return variables

关键是那个 _with_selectorcore/workflow/system_variables.py:102)——它保证每个变量的地址第一段是对的前缀:

# core/workflow/system_variables.py:102
def _with_selector(variable: Variable, node_id: str) -> Variable:
    selector = [node_id, variable.name]
    if list(variable.selector) == selector:
        return variable                              # 已经对了:原样返回,不复制
    return variable.model_copy(update={"selector": selector})   # 否则:改地址再复制一份
三个 *(...)把 system / environment / conversation 三类变量分别"贴标签"(盖前缀)后铺平进同一个列表。星号是 Python 的解包:把生成器里的元素摊进外层 list。
_with_selector 的短路★如果地址已经正确就直接返回原对象,只有不对时才 model_copy 复制一份改地址。少造一次对象——热路径上的小优化。
rag 按 belong_to_node_id 分组RAG 流水线变量归属于具体节点,所以先按节点分组聚成一个 dict,再整体作为一个变量存到 ("rag", node_id) 下。
💡 设计取舍:为什么变量要"不可变 + copy"?Variable 是 pydantic 模型,改字段要靠 model_copy 造新对象而不是原地改。好处是同一个变量对象在多处被引用时,不会因为某处改地址而"祸及"别处(不可变数据的经典优势)。代价是可能多造几个对象——所以才有 _with_selector 里"已正确就不复制"的短路。用不可变换安全,用短路省开销
L06

节点跑完,输出怎么写回池

开机装的是初始料。节点跑起来后产出的新数据,靠 variable_pool_initializer.py 里的两个小函数写回池(core/workflow/variable_pool_initializer.py:8):

# core/workflow/variable_pool_initializer.py:8
def add_variables_to_pool(variable_pool, variables) -> None:
    for variable in variables:
        variable_pool.add(variable.selector, variable)      # 逐个按自带地址塞进池

def add_node_inputs_to_pool(variable_pool, *, node_id, inputs, aliases=()):
    """Store node inputs under the primary node id and any compatible aliases."""
    node_ids = [node_id]
    for alias in aliases:                                   # ① 主 id + 若干别名
        if alias not in node_ids:
            node_ids.append(alias)
    for current_node_id in node_ids:
        for key, value in inputs.items():
            variable_pool.add((current_node_id, key), value)   # ② 每个字段都存一份
add_variables_to_pool最朴素的写法:变量自己知道自己的地址(variable.selector),挨个 add 进去即可。
add_node_inputs_to_pool 的 aliases★同一份输入,除了存到主 node_id 下,还能存到若干"别名"下。为什么?因为节点可能改过 id(或有历史 id),旧画布里的 {{#旧id.字段#}} 引用还得能读到——别名就是"同一个抽屉挂多个门牌",保证老引用不失效
(current_node_id, key)输入是个字典,逐字段拆开,每个字段拼一个 (节点id, 字段名) 地址存进去。这样 {{#节点id.字段#}} 就能精确读到某个字段。
变量池 = 一块按前缀分区的共享白板 VariablePool(字典:selector → Segment) sys query / files user_id … env app 级常量 conversation 跨轮会话状态 节点 id llm.text … 开机:build_bootstrap_variables 装 sys/env/conversation/rag 运行中:add_node_inputs_to_pool 把每个节点的输出写到"节点 id"抽屉 读:{{#sys.query#}} → 拆成 ("sys","query") → pool.get() → 真实值 节点彼此不认识,只认这块白板 —— 这就是解耦
图注:四类前缀 + 节点 id 各占一片区域;写靠 add,读靠 get,selector 就是地址。
L07

读取变量 + 今日小结

读取端也有一组便捷函数(core/workflow/system_variables.py:140),把"查池 + 拆 Segment"包成一步:

# core/workflow/system_variables.py:140
def get_system_segment(variable_pool, key) -> Segment | None:
    return variable_pool.get(system_variable_selector(key))   # 按 ("sys",key) 查

def get_system_value(variable_pool, key) -> Any:
    segment = get_system_segment(variable_pool, key)
    return None if segment is None else segment.value         # 拆出裸值

def get_all_system_variables(variable_pool):
    return variable_pool.get_by_prefix(SYSTEM_VARIABLE_NODE_ID)   # 一次拿回整个 sys 抽屉
📝 真实值:一句 {{#sys.query#}} 的完整旅程 用户在多轮对话里问「帮我总结这份合同」并上传了 contract.pdf → 引擎开机时 build_system_variables(query="帮我总结这份合同", files=[contract.pdf], conversation_id="c-88", user_id="u-3", dialogue_count=2) → 经 _normalize 清洗(files 已非空、无 None)→ 造出 5 个 Variable,地址分别是 ("sys","query")("sys","files")… → build_bootstrap_variables 装进池 → LLM 节点渲染 Prompt 时遇到 {{#sys.query#}} → 拆成 ("sys","query")get_system_value 查回 "帮我总结这份合同" → 替换进 Prompt → 发给模型。你写的那半行 {{#sys.query#}},背后是"造变量→装池→按地址取值"这一整条链。

👶 小白:selector、Segment、Variable 这几个名字好像,谁是谁?

👨‍🏫 老师:一句话分清——selector 是"地址"(如 ("sys","query"),就是字典的 key);Segment 是"带类型的值"(把裸字符串/列表包一层,知道自己是文本还是数组);Variable 是"地址 + 值 + 名字的完整条目"selector + Segment + name,是真正存进池的东西)。类比一封信:selector 是收件地址,Segment 是信的内容,Variable 是贴好地址、装好内容的整个信封。

🧠 今天你应该能回答

  • 工作流节点之间靠什么传数据?(全局变量池 VariablePool,一块共享白板)
  • 一个变量的地址长什么样?(selector,形如 (命名空间/节点id, 字段名)
  • 四个命名空间前缀是哪些?(sys/env/conversation/rag,外加节点 id)
  • 系统变量有哪些、由谁造?(SystemVariableKey 枚举钉死,build_system_variables 造)
  • 开机装料靠谁?(build_bootstrap_variables,给四类变量盖前缀)
  • 节点输出怎么写回、老引用为什么不失效?(add_node_inputs_to_pool + aliases 多门牌)

✋ 10 分钟动手

cd /Users/bitmart/work/codes/github/AI_WORK/dify/api

# 1. 四个前缀(就 4 行)
cat core/workflow/variable_prefixes.py

# 2. 系统变量清单 + 构造
sed -n '22,39p'   core/workflow/system_variables.py    # SystemVariableKey 枚举
sed -n '81,96p'   core/workflow/system_variables.py    # build_system_variables
sed -n '109,138p' core/workflow/system_variables.py    # build_bootstrap_variables

# 3. 写回池
cat core/workflow/variable_pool_initializer.py         # 全文就 28 行

# 4. 数一数系统变量有几个
grep -c '= "' core/workflow/system_variables.py
明日预告 · Day 12:今天变量池把"数据"存好了,可工作流跑起来是动态的——节点开始了、模型吐字了、节点成功了…这些事件怎么一条条推给前端做成打字机效果?明天我们进 task_pipeline,看工作流的事件系统与流式输出:几十种 Queue 事件如何被分发、如何转成 SSE 流。
← Day 10 Day 12 · 工作流事件与流式输出 →