工作流编排:画布背后的那套内核

执行顺序不是你画的位置决定的,是依赖关系算出来的——看懂这句话,换任何编排平台都不用重学。

30″30 秒看懂工作流编排

把编排想成一条工厂流水线:每个工位只干一件事,传送带把半成品往下送,还有一个料箱跟着传送带一起走——每个工位从箱里取自己要的料,加工完把成品放回箱里。

分拣工位先看这单是什么类型,取料工位去查知识库和订单系统(这两个谁也不等谁,可以同时开工),组装工位把料拼成提示词,加工工位交给大模型生成。最后箱子里那件成品,就是给用户的回复。

工作流编排 = 一条流水线:工位、传送带、料箱 1 分拣工位 判断这单是什么类型 2 取料工位 查知识库 2 取料工位 查订单系统 3 组装工位 把料拼成提示词 4 加工工位 交给大模型生成 成品 回复 这两个工位谁也不等谁 → 可以并行 随传送带一起走的那个料箱 = 变量表(state) 每个工位从箱里取自己要的料,加工完把成品放回箱里。箱子是共享的、可被覆盖的、没有类型检查的—— 这三个特点凑在一起,就是编排里最难查的三类 bug 的全部来源。 question → intent → docs / order → prompt → answer 箱子里的料越走越多,只增不改
图① 30 秒看懂:工位、传送带、跟着走的料箱
流水线上的东西对应的技术概念它到底是什么
工位节点(node)读几个变量、干一件事、写一个变量。大模型节点和代码节点结构完全一样,只是中间干的事不同
传送带连线(edge)其实不用手画——一个节点读了谁写的变量,边就自动存在了
跟着走的料箱变量表(state)所有节点共用的一份字典,只增不改。编排里最难查的 bug 全出在这个箱子上
两个工位同时开工并行分层能不能并行由依赖关系决定,不是由你在画布上画的位置决定
半成品退回上一站回边 / 循环节点让流程能「改到满意为止」,代价是不再保证一定会停
某个工位断料了节点失败先判断该不该重试,再决定这一路缺了算不算数
⛔ 整页只有一条铁律 执行顺序不是你画的位置决定的,是依赖关系算出来的。画布上从左到右排得再整齐,只要下游没读上游的变量,它们就是并行的;只要读了,画得再远也是串行的。看不懂这句话,就会一直在画布上调位置,却怎么也调不对顺序。

这一页会用 120 行代码把这套内核写出来跑给你看。看懂之后,换到任何一个编排平台——画布长什么样、节点叫什么名字——你都能立刻对上号。

01概念:编排到底在编排什么

节点的三段结构、连线为什么不用手画、以及对话式与流程式的分界

1.1 一个节点只有三段

不管平台把它叫「大模型节点」「代码节点」「知识库检索节点」还是「HTTP 请求节点」,剥开之后都是同一副骨架:

节点解剖:不管平台把它叫什么,结构都是这三段 「大模型节点」「代码节点」「知识库节点」的差别只在中间那段干了什么 一个节点 1 reads —— 声明我要读哪些变量 reads = ["docs", "order", "intent"] 2 fn —— 真正干活的那段逻辑 调模型 / 跑代码 / 查知识库 / 请求接口,节点类型的差别只在这里 3 writes —— 声明我把结果写进哪个变量 writes = "prompt" 连线不用手画 一个节点读了谁写的变量, 边就自动存在了。 retrieve --docs--> compose lookup_order --order--> compose compose --prompt--> generate 所以「连线漏了」的本质不是少画一根线, 是下游读了一个没人写过的变量名。 所有节点共用的那一份 state question intent docs / order prompt answer 只增不改
图② 节点解剖:reads / fn / writes 三段,连线自动推出
写法作用
reads["docs", "order", "intent"]声明我要读哪些变量。只有声明过的才交给它,防止节点偷偷依赖没声明的东西
fn一段函数真正干活:调模型、跑代码、查库、请求接口。节点类型的差别只在这里
writes"prompt"声明结果写进哪个变量。一个节点只写一个变量,多写就容易撞名
为什么强调「只写一个变量」 因为变量表是全局共享、可被覆盖、没有类型检查的。一个节点写多处,等于在全局命名空间里多埋几颗雷。后面 2.4 会看到同名覆盖翻车的完整过程。

1.2 连线其实不用手画

这是理解编排最关键的一个转折:你在画布上拖的那根线,只是依赖关系的可视化,不是依赖关系本身。

真正的依赖来自变量:compose 读了 docs,而 docsretrieve 写的,那么 retrieve → compose 这条边就已经存在了,不管你有没有把它画出来。所以引擎可以自动推出全部连线:

推出的边依据含义
retrieve --docs--> composecompose 读 docscompose 必须等 retrieve 跑完
lookup_order --order--> composecompose 读 order同上
classify --intent--> lookup_orderlookup_order 读 intent要先知道意图才去查订单
compose --prompt--> generategenerate 读 prompt提示词拼好才能生成

这个视角有个很实用的推论:「连线漏了」的本质不是少画一根线,而是下游读了一个没人写过的变量名。所以排查断链时,别盯着画布找哪根线没连上,去查变量名对不对得上——这快得多。

1.3 对话式与流程式:两种入口,一个内核

多数平台会提供两种编排形态,名字各不相同,但分界线是一致的:

形态入口是什么有没有会话态适合什么
对话式用户的每一句话有,多轮上下文自动带客服、助手、问答——用户会追问的场景
流程式一次输入 / 一批数据没有,每次都是干净的批量处理、定时任务、被别的系统调用

两者内核完全相同——都是节点、变量表、拓扑执行。区别只在于:对话式在变量表里预置了会话历史,并且默认把最后一个节点的输出当作回复流式吐出去。

⚠️ 选错形态的代价 用流程式去做需要追问的客服,结果是每轮都失忆;用对话式去做批量处理,结果是会话历史越滚越长,提示词成本随轮次线性上涨,并且互不相关的两条数据会互相污染上下文。选型时先问一句:用户会不会追问?

1.4 什么时候不该用编排

编排擅长表达有向、分支有限、每步职责清晰的流程。它不擅长:动态生成的流程结构、精细的并发与重试控制、深度嵌套的状态机。

判断信号很明确——当你开始用一堆条件节点模拟 if-else 嵌套,或者不停往代码节点里塞逻辑时,画布就已经变成「用鼠标写代码」了。这时候直接写代码更快、更好维护,也更容易测试。

02原理:拓扑排序、并行分层、以及变量表上的三类故障

执行顺序怎么算出来、哪些节点真能并行、以及为什么最难查的 bug 都出在同一个地方

2.1 执行顺序是算出来的

引擎拿到一堆节点后,做的第一件事是拓扑排序:统计每个节点的入度(它依赖几个上游),把入度为 0 的先放进队列,跑完一个就给它的下游入度减一,减到 0 就入队。

这个过程顺带解决了一个重要问题——环检测。如果排完之后还有节点没被排进去,说明它们的入度永远减不到 0,也就是画布上存在回边。

执行顺序不是你画的位置决定的,是依赖关系算出来的 按层分组:同一层互不依赖,可以并行 第1层 第2层 第3层 第4层 classify retrieve 这两个可并行 lookup_order 要等 classify 给出 intent compose 要等 docs + order + intent 三样齐 generate 要等 prompt 出现回边,排序就失败 classify compose generate answer 回流 入度永远减不到 0,节点轮不到执行 实跑输出: 执行顺序:classify → retrieve → lookup_order → compose → generate 造环后捕获:工作流里有环,这些节点永远轮不到执行:['classify','compose','generate','lookup_order']
图③ 拓扑排序与并行分层,以及回边导致的排序失败
120 行的工作流引擎:节点、边、拓扑排序、变量表
"""一个 120 行的工作流引擎:把画布上的节点和连线跑起来。

所有可视化编排平台的内核都是这套东西——节点、边、拓扑排序、
一份在节点之间传递的变量表。看懂这个文件,你就看懂了画布背后
到底发生了什么,换任何平台都不用重学。

纯标准库,直接 python3 dag_engine.py 就能跑。
"""

from collections import deque


class Node:
    """一个节点 = 一个名字 + 一个函数 + 声明它要读哪些变量、写哪个变量。

    平台上那些「开始节点」「大模型节点」「代码节点」,
    差别只在 fn 里干了什么,结构上完全一样。
    """

    def __init__(self, name, fn, reads, writes):
        self.name = name
        self.fn = fn
        self.reads = reads      # 依赖的变量名列表
        self.writes = writes    # 产出的变量名

    def run(self, state):
        # 只把声明过的变量交给节点,防止节点偷偷依赖没声明的东西
        kwargs = {k: state[k] for k in self.reads}
        return self.fn(**kwargs)


class Workflow:
    def __init__(self):
        self.nodes = {}
        self.producer = {}   # 变量名 -> 产出它的节点名

    def add(self, node):
        if node.writes in self.producer:
            raise ValueError(
                "变量 %s 被两个节点写入:%s%s —— "
                "画布上最难查的 bug 就是这种覆盖"
                % (node.writes, self.producer[node.writes], node.name))
        self.nodes[node.name] = node
        self.producer[node.writes] = node.name
        return self

    def _edges(self):
        """边不用手画:一个节点读了谁写的变量,边就自动存在。"""
        deps = {name: set() for name in self.nodes}
        for name, node in self.nodes.items():
            for var in node.reads:
                if var in self.producer:
                    deps[name].add(self.producer[var])
        return deps

    def topo_order(self):
        """拓扑排序:算出节点的执行顺序,顺便检测环。"""
        deps = self._edges()
        indeg = {n: len(d) for n, d in deps.items()}
        # 反向表:某节点跑完后,谁的入度可以减一
        downstream = {n: [] for n in self.nodes}
        for n, d in deps.items():
            for parent in d:
                downstream[parent].append(n)

        q = deque(sorted(n for n, k in indeg.items() if k == 0))
        order = []
        while q:
            n = q.popleft()
            order.append(n)
            for child in sorted(downstream[n]):
                indeg[child] -= 1
                if indeg[child] == 0:
                    q.append(child)

        if len(order) != len(self.nodes):
            stuck = sorted(set(self.nodes) - set(order))
            raise ValueError(
                "工作流里有环,这些节点永远轮不到执行:%s" % stuck)
        return order

    def parallel_layers(self):
        """按层分组:同一层的节点互不依赖,可以并行跑。"""
        deps = self._edges()
        done, layers = set(), []
        while len(done) < len(self.nodes):
            layer = sorted(
                n for n in self.nodes
                if n not in done and deps[n] <= done)
            if not layer:
                raise ValueError("工作流里有环,无法分层")
            layers.append(layer)
            done |= set(layer)
        return layers

    def run(self, initial_state, verbose=True):
        state = dict(initial_state)
        for name in self.topo_order():
            node = self.nodes[name]
            missing = [k for k in node.reads if k not in state]
            if missing:
                raise KeyError(
                    "节点 %s 读不到变量 %s —— "
                    "多半是连线漏了,或者上游节点改了输出名"
                    % (name, missing))
            value = node.run(state)
            state[node.writes] = value
            if verbose:
                print("  [%-10s] %s = %r" % (name, node.writes, value))
        return state


def build_demo():
    """一个真实形状的流程:分类 → 两路并行取数 → 汇总 → 生成。"""
    wf = Workflow()
    wf.add(Node("classify", lambda question: (
        "售后" if "退" in question or "坏" in question else "咨询"),
        reads=["question"], writes="intent"))
    wf.add(Node("retrieve", lambda question: (
        ["文档片段A", "文档片段B"]),
        reads=["question"], writes="docs"))
    wf.add(Node("lookup_order", lambda question, intent: (
        {"order_id": "SO-2026-0918", "status": "已发货"}
        if intent == "售后" else None),
        reads=["question", "intent"], writes="order"))
    wf.add(Node("compose", lambda docs, order, intent: (
        "意图=%s | 依据=%d 条 | 订单=%s"
        % (intent, len(docs), order["status"] if order else "无")),
        reads=["docs", "order", "intent"], writes="prompt"))
    wf.add(Node("generate", lambda prompt: "【回复】" + prompt,
        reads=["prompt"], writes="answer"))
    return wf


def main():
    wf = build_demo()

    print("执行顺序(拓扑排序):")
    print("  " + " → ".join(wf.topo_order()))

    print("\n可并行的层:")
    for i, layer in enumerate(wf.parallel_layers(), 1):
        tag = "(可并行)" if len(layer) > 1 else ""
        print("  第%d层: %s %s" % (i, "、".join(layer), tag))

    print("\n执行过程:")
    final = wf.run({"question": "我买的杯子坏了,能退吗"})
    print("\n最终输出:%s" % final["answer"])

    # 环检测:把 generate 的输出接回 classify 的输入
    print("\n--- 故意造一个环 ---")
    bad = build_demo()
    bad.nodes["classify"].reads = ["question", "answer"]
    try:
        bad.topo_order()
    except ValueError as e:
        print("  捕获:%s" % e)


if __name__ == "__main__":
    main()

实跑输出:

执行顺序、并行分层、以及故意造环后的报错
执行顺序(拓扑排序):
  classify → retrieve → lookup_order → compose → generate

可并行的层:
  第1层: classify、retrieve (可并行)
  第2层: lookup_order 
  第3层: compose 
  第4层: generate 

执行过程:
  [classify  ] intent = '售后'
  [retrieve  ] docs = ['文档片段A', '文档片段B']
  [lookup_order] order = {'order_id': 'SO-2026-0918', 'status': '已发货'}
  [compose   ] prompt = '意图=售后 | 依据=2 条 | 订单=已发货'
  [generate  ] answer = '【回复】意图=售后 | 依据=2 条 | 订单=已发货'

最终输出:【回复】意图=售后 | 依据=2 条 | 订单=已发货

--- 故意造一个环 ---
  捕获:工作流里有环,这些节点永远轮不到执行:['classify', 'compose', 'generate', 'lookup_order']
注意 add() 里那个检查 两个节点写同一个变量时直接抛错,并在错误信息里写明是哪两个节点。画布上最难查的 bug 就是这种静默覆盖——它不报错,只是后写的把先写的盖掉,下游拿到的数据莫名其妙。在引擎层面禁掉它,比事后排查划算得多。

2.2 并行分层:谁真的能同时跑

把节点按「依赖是否已全部满足」分层:第一层是不依赖任何人的,第二层是只依赖第一层的,以此类推。同一层内的节点互不依赖,可以并行

示例流程的分层结果是:

节点能否并行
第 1 层classify、retrieve可并行——两个都只读 question,谁也不等谁
第 2 层lookup_order单节点,要等 classify 给出 intent
第 3 层compose单节点,要等 docs + order + intent 三样齐
第 4 层generate单节点,要等 prompt

这里有个反直觉的点:retrieve 在画布上可能被画在很靠后的位置,但它其实第一层就能开跑。因为它只读 question,不依赖任何中间结果。很多人画布画得很整齐,却没意识到自己把两个本可并行的节点串成了一条线——只要给后面那个节点加一个它根本不需要的输入,并行就没了。

2.3 变量表:便利与风险来自同一个设计

节点之间靠一份共享的变量表通信。这份表有三个特点,它们既是编排好用的原因,也是它最容易出事的原因:

1全局共享

任何节点都能读任何已有变量,不用层层传参。代价是命名空间只有一个,撞名就覆盖。

2可被覆盖

写入就是赋值,不会提示「这个名字已经有人用了」。代价是静默覆盖

3没有类型检查

今天写进去是列表,明天换个节点写成字符串,引擎不管。代价是类型漂移

4只增不改

这一条是好事:变量一旦写入就不再变动,排障时可以完整回放每一步的中间值。

2.4 三类故障,全部来自变量表

带写入追踪的变量表:把三类隐患捞出来
"""变量作用域:画布上最难查的那类 bug,根源都在这里。

节点之间靠一份共享的变量表通信。这份表是全局的、可被覆盖的、
没有类型检查的 —— 三个特点凑在一起,就构成了编排里最常见的
三种故障:改名断链、同名覆盖、类型漂移。

纯标准库,直接 python3 state_scope.py 就能跑。
"""


class Scope:
    """带写入追踪的变量表:谁写的、写了几次、被谁读过。"""

    def __init__(self):
        self.values = {}
        self.written_by = {}     # 变量 -> [写它的节点, ...]
        self.read_by = {}        # 变量 -> [读它的节点, ...]

    def write(self, node, key, value):
        self.values[key] = value
        self.written_by.setdefault(key, []).append(node)

    def read(self, node, key):
        if key not in self.values:
            raise KeyError(
                "节点 %s 读不到 %s。三种可能:上游改了输出名、"
                "连线漏了、或者这个节点被排到了上游前面" % (node, key))
        self.read_by.setdefault(key, []).append(node)
        return self.values[key]

    def audit(self):
        """跑完之后做一次体检,把三类隐患捞出来。"""
        problems = []

        for key, writers in self.written_by.items():
            if len(writers) > 1:
                problems.append(
                    ("同名覆盖", key,
                     "被 %s 先后写入,后写的把先写的盖掉了" % "、".join(writers)))

        for key in self.written_by:
            if key not in self.read_by:
                problems.append(
                    ("写了没人读", key,
                     "产出后无人使用,要么是废节点,要么是下游连线漏了"))

        return problems


def demo_type_drift():
    """类型漂移:同一个变量在不同分支里是不同类型。"""
    sc = Scope()

    # A 分支:检索节点返回列表
    sc.write("retrieve", "docs", ["片段1", "片段2"])
    n = len(sc.read("compose", "docs"))
    print("  A 分支:docs 是 list,len=%d,一切正常" % n)

    # B 分支:换了个节点,它返回的是拼好的字符串
    sc.write("retrieve_v2", "docs", "片段1\n片段2")
    v = sc.read("compose", "docs")
    print("  B 分支:docs 变成 str,len=%d —— 数字还在,含义已经错了"
          % len(v))
    print("  ⚠ 这类 bug 不报错,只是答案悄悄变差,最难查")


def demo_rename_break():
    """改名断链:上游改了输出名,下游还在读旧名字。"""
    sc = Scope()
    sc.write("retrieve", "chunks", ["片段1"])   # 原来叫 docs,改成了 chunks
    try:
        sc.read("compose", "docs")
    except KeyError as e:
        print("  捕获:%s" % e)


def demo_overwrite():
    """同名覆盖:两个节点都往同一个变量里写。"""
    sc = Scope()
    sc.write("lookup_order", "result", {"status": "已发货"})
    sc.write("lookup_ticket", "result", {"ticket": "T-001"})
    got = sc.read("compose", "result")
    print("  compose 拿到的是 %r —— 订单信息已经被工单覆盖掉了" % got)
    sc.write("dead_node", "unused", 42)
    return sc


def main():
    print("【故障一】类型漂移")
    demo_type_drift()

    print("\n【故障二】改名断链")
    demo_rename_break()

    print("\n【故障三】同名覆盖")
    sc = demo_overwrite()

    print("\n作用域体检:")
    for kind, key, why in sc.audit():
        print("  ❌ %-8s %-10s %s" % (kind, key, why))

    print("\n三条防线:")
    print("  1. 变量名带节点前缀,如 retrieve_docs,天然不会撞名")
    print("  2. 每个节点的输出只写一个变量,别让一个节点写多处")
    print("  3. 上线前跑一次体检,把「写了没人读」的节点清掉")


if __name__ == "__main__":
    main()

实跑,三类故障依次复现:

类型漂移、改名断链、同名覆盖
【故障一】类型漂移
  A 分支:docs 是 list,len=2,一切正常
  B 分支:docs 变成 str,len=7 —— 数字还在,含义已经错了
  ⚠ 这类 bug 不报错,只是答案悄悄变差,最难查

【故障二】改名断链
  捕获:'节点 compose 读不到 docs。三种可能:上游改了输出名、连线漏了、或者这个节点被排到了上游前面'

【故障三】同名覆盖
  compose 拿到的是 {'ticket': 'T-001'} —— 订单信息已经被工单覆盖掉了

作用域体检:
  ❌ 同名覆盖     result     被 lookup_order、lookup_ticket 先后写入,后写的把先写的盖掉了
  ❌ 写了没人读    unused     产出后无人使用,要么是废节点,要么是下游连线漏了

三条防线:
  1. 变量名带节点前缀,如 retrieve_docs,天然不会撞名
  2. 每个节点的输出只写一个变量,别让一个节点写多处
  3. 上线前跑一次体检,把「写了没人读」的节点清掉
故障表现为什么难查
类型漂移docs 从 list 变成 str,len() 从 2 变成 7不报错。数字还在,含义已经错了,只是答案悄悄变差
改名断链上游改叫 chunks,下游还读 docs报错但指向不明——可能是改名、漏连线、或顺序排错
同名覆盖订单信息被工单信息盖掉不报错。下游拿到的是另一个节点的产出,看上去像模型答错了
⛔ 三条防线变量名带节点前缀,如 retrieve_docs,天然不会撞名;② 每个节点只写一个变量;③ 上线前跑一次作用域体检,把「写了没人读」的节点清掉——那种节点要么是废的,要么是下游连线漏了。

03最小可跑:失败、并行、循环三件事

流程在顺利时都长得一样,区别全在出事的时候

前面那个引擎有个隐含假设:每个节点都会成功返回。真实流程里这个假设一次都不成立——模型会超时,接口会限流,网络会抖。这一节把三件真实会发生的事补上。

3.1 节点失败:先判断该不该重试

很多人给所有节点统一开「重试 3 次」,这是错的。重试只对「再试一次可能就好了」的错误有意义,对参数写错、鉴权失败、内容被拦这类错误,重试只会把一次失败放大成四次。

节点失败了:先问「该不该重试」,再问「重试怎么排队」,最后问「彻底失败算不算数」 1 该不该重试 ✓ 再试一次可能就好了 timeout · rate_limited upstream_5xx ✗ 再试一百次也一样 bad_request · auth_failed content_filtered 2 重试怎么排队 不带抖动:所有请求同一刻齐射 [0.4, 0.8, 1.6, 3.2, 6.4] 带抖动:把重试摊开到一段区间里 [0.38, 0.75, 1.29, 3.2, 6.13] 不加抖动 = 把刚缓过来的下游再打死一次 3 失败了算不算数 关键路失败 → 整条流程停 缺了就答不了的,比如订单状态 宁可报错,也别编一个答案 可选路失败 → 标注后继续 降级值要是业务能接受的保守答案 并在提示词里写明「本次未取到」 三种真实走向 抖动两次后恢复 第1次 timeout → 等 0.27s 第2次 timeout → 等 0.72s 第3次 ✅ 成功 重试机制恰好发挥作用的样子 一直超时,走降级 4 次全 timeout,重试用尽 → 降级为「本次未取到外部资料」 流程没挂,但下游必须知道 这一路是空的 参数写错,重试毫无意义 第1次 bad_request → 不可重试,立即放弃 对这类错误重试,只会把 一次失败放大成四次
图④ 失败处理三件套:可重试判定、退避抖动、关键路与降级
重试、退避、降级三件套
"""节点失败了怎么办:重试、超时、降级三件套。

画布上每个调外部服务的节点都会失败 —— 模型超时、接口限流、
网络抖动。平台一般在节点设置里提供这三个开关,但开关背后的
策略必须你自己想清楚,否则要么雪上加霜,要么整条流程挂掉。

纯标准库,直接 python3 retry_fallback.py 就能跑。
"""

import random
import time

random.seed(20260918)   # 固定种子,保证输出可复现


class NodeFailed(Exception):
    pass


def backoff_delays(max_retries, base=0.4, cap=8.0, jitter=True):
    """指数退避 + 抖动,返回每次重试前要等的秒数。

    不加抖动时,被同一次故障打挂的所有请求会在同一时刻一起重试,
    把刚缓过来的下游再次打死 —— 这叫重试风暴。
    """
    delays = []
    for i in range(max_retries):
        d = min(base * (2 ** i), cap)
        if jitter:
            d = random.uniform(d * 0.5, d)
        delays.append(round(d, 2))
    return delays


# 哪些错误值得重试:只有「再试一次可能就好了」的才重试
RETRYABLE = {"timeout", "rate_limited", "upstream_5xx"}
NON_RETRYABLE = {"bad_request", "auth_failed", "content_filtered"}


def call_with_retry(node_name, fn, max_retries=3, sleep=False):
    """带重试的节点调用。返回 (成功?, 结果或错误, 尝试次数)。"""
    delays = backoff_delays(max_retries)
    for attempt in range(max_retries + 1):
        try:
            return True, fn(attempt), attempt + 1
        except NodeFailed as e:
            kind = str(e)
            if kind in NON_RETRYABLE:
                print("    第%d次失败(%s):不可重试,立即放弃"
                      % (attempt + 1, kind))
                return False, kind, attempt + 1
            if attempt == max_retries:
                print("    第%d次失败(%s):重试用尽" % (attempt + 1, kind))
                return False, kind, attempt + 1
            d = delays[attempt]
            print("    第%d次失败(%s):等待 %.2fs 后重试"
                  % (attempt + 1, kind, d))
            if sleep:
                time.sleep(d)
    return False, "unknown", max_retries + 1


def make_flaky(fail_times, kind="timeout"):
    """造一个前 N 次必失败、之后成功的假节点。"""
    def fn(attempt):
        if attempt < fail_times:
            raise NodeFailed(kind)
        return "第%d次成功拿到结果" % (attempt + 1)
    return fn


def run_with_fallback(node_name, primary, fallback_value, max_retries=3):
    """主路径失败后走降级:宁可给一个保守答案,也别整条流程挂掉。"""
    print("  节点 %s:" % node_name)
    ok, result, tries = call_with_retry(node_name, primary, max_retries)
    if ok:
        print("    ✅ %s(共 %d 次尝试)" % (result, tries))
        return result
    print("    ⚠ 主路径失败(%s),降级为:%r" % (result, fallback_value))
    return fallback_value


def main():
    print("退避序列(base=0.4s,上限 8s,带抖动):")
    print("  %s" % backoff_delays(5))
    print("  不带抖动:%s  <- 所有请求会在同一刻齐射"
          % backoff_delays(5, jitter=False))

    print("\n【场景一】抖动两次后恢复")
    run_with_fallback("call_model", make_flaky(2), fallback_value=None)

    print("\n【场景二】一直超时,走降级")
    run_with_fallback("search_web", make_flaky(99),
                      fallback_value="(本次未取到外部资料)")

    print("\n【场景三】参数写错了,重试毫无意义")
    run_with_fallback("call_api", make_flaky(99, "bad_request"),
                      fallback_value="(接口调用被跳过)")

    print("\n三条判断依据:")
    print("  • 可重试:%s" % "、".join(sorted(RETRYABLE)))
    print("  • 不可重试:%s" % "、".join(sorted(NON_RETRYABLE)))
    print("  • 降级值必须是业务能接受的保守答案,不是随便填个空字符串")


if __name__ == "__main__":
    main()
三种真实走向的实跑结果
退避序列(base=0.4s,上限 8s,带抖动):
  [0.38, 0.75, 1.29, 3.2, 6.13]
  不带抖动:[0.4, 0.8, 1.6, 3.2, 6.4]  <- 所有请求会在同一刻齐射

【场景一】抖动两次后恢复
  节点 call_model:
    第1次失败(timeout):等待 0.27s 后重试
    第2次失败(timeout):等待 0.72s 后重试
    ✅ 第3次成功拿到结果(共 3 次尝试)

【场景二】一直超时,走降级
  节点 search_web:
    第1次失败(timeout):等待 0.23s 后重试
    第2次失败(timeout):等待 0.73s 后重试
    第3次失败(timeout):等待 0.89s 后重试
    第4次失败(timeout):重试用尽
    ⚠ 主路径失败(timeout),降级为:'(本次未取到外部资料)'

【场景三】参数写错了,重试毫无意义
  节点 call_api:
    第1次失败(bad_request):不可重试,立即放弃
    ⚠ 主路径失败(bad_request),降级为:'(接口调用被跳过)'

三条判断依据:
  • 可重试:rate_limited、timeout、upstream_5xx
  • 不可重试:auth_failed、bad_request、content_filtered
  • 降级值必须是业务能接受的保守答案,不是随便填个空字符串
错误类型该不该重试理由
timeout下游可能只是慢了一下
rate_limited等一会儿配额就回来了
upstream_5xx对方服务端临时故障
bad_request不该参数本身就是错的,试一百次还是错
auth_failed不该凭证问题,重试解决不了
content_filtered不该内容被拦,同样的输入结果相同
⚠️ 退避一定要加抖动 看实跑输出里那两行对比:不带抖动是 [0.4, 0.8, 1.6, 3.2, 6.4],带抖动是 [0.38, 0.75, 1.29, 3.2, 6.13]。差别看着不大,但不带抖动时,被同一次故障打挂的所有请求会在同一时刻一起重试,把刚缓过来的下游再次打死。这叫重试风暴,是把小故障放大成大事故的经典路径。

3.2 并行:难点在收口,不在扇出

扇出很简单,几行线程池就完了。真正要想清楚的是:三路里有一路失败或超时,这次算成功还是失败?

并行扇出与两种收口策略
"""并行与聚合:扇出去、收回来,以及收不齐时怎么办。

画布上两个节点并排画着,不等于它们真的同时跑。真正决定能不能
并行的是依赖关系 —— 谁也不读谁的输出,才谈得上并行。

难点从来不在「怎么扇出去」,而在「收回来时有一路失败/超时」。
纯标准库(concurrent.futures),直接 python3 parallel_gather.py 就能跑。
"""

import random
import time
from concurrent.futures import ThreadPoolExecutor, as_completed

random.seed(918)


def branch(name, cost, fail=False):
    """造一个耗时 cost 秒的分支任务。"""
    def fn():
        time.sleep(cost)
        if fail:
            raise RuntimeError("%s 调用失败" % name)
        return "%s 的结果" % name
    fn.__name__ = name
    return fn


def gather_all_or_nothing(tasks, timeout):
    """严格模式:任一路失败,整体失败。

    适合「缺了就不能答」的场景,比如必须拿到订单状态才能回复。
    """
    results, errors = {}, {}
    with ThreadPoolExecutor(max_workers=len(tasks)) as pool:
        futs = {pool.submit(fn): name for name, fn in tasks.items()}
        try:
            for fut in as_completed(futs, timeout=timeout):
                name = futs[fut]
                try:
                    results[name] = fut.result()
                except Exception as e:
                    errors[name] = str(e)
        except TimeoutError:
            for fut, name in futs.items():
                if name not in results and name not in errors:
                    errors[name] = "超时未返回"
    ok = not errors
    return ok, results, errors


def gather_best_effort(tasks, timeout, required):
    """尽力模式:关键路必须成功,其余缺了就标注为空。

    适合「多一份资料更好,少一份也能答」的场景。
    """
    ok, results, errors = gather_all_or_nothing(tasks, timeout)
    missing_required = [k for k in required if k not in results]
    usable = not missing_required
    return usable, results, errors, missing_required


def show(title, usable, results, errors, missing=None):
    print("\n%s】" % title)
    for k, v in sorted(results.items()):
        print("  ✅ %-12s %s" % (k, v))
    for k, v in sorted(errors.items()):
        print("  ❌ %-12s %s" % (k, v))
    if missing:
        print("  ⛔ 关键路缺失:%s" % "、".join(missing))
    print("  → 本次%s继续往下走" % ("可以" if usable else "不能"))


def main():
    # 三路并行:检索最快,订单中等,网页搜索最慢
    tasks = {
        "retrieve": branch("retrieve", 0.20),
        "order": branch("order", 0.35),
        "websearch": branch("websearch", 0.55),
    }

    t0 = time.time()
    ok, res, err = gather_all_or_nothing(tasks, timeout=2.0)
    elapsed = time.time() - t0
    show("三路全成功(严格模式)", ok, res, err)
    print("  耗时 %.2fs —— 串行要 %.2fs,并行只花最慢那一路的时间"
          % (elapsed, 0.20 + 0.35 + 0.55))

    # 网页搜索挂了
    tasks2 = dict(tasks)
    tasks2["websearch"] = branch("websearch", 0.30, fail=True)

    ok, res, err = gather_all_or_nothing(tasks2, timeout=2.0)
    show("一路失败(严格模式)", ok, res, err)

    usable, res, err, missing = gather_best_effort(
        tasks2, timeout=2.0, required=["retrieve", "order"])
    show("一路失败(尽力模式,关键路=retrieve+order)",
         usable, res, err, missing)

    # 关键路自己挂了
    tasks3 = dict(tasks)
    tasks3["order"] = branch("order", 0.20, fail=True)
    usable, res, err, missing = gather_best_effort(
        tasks3, timeout=2.0, required=["retrieve", "order"])
    show("关键路失败(尽力模式)", usable, res, err, missing)

    print("\n收口的三个决定:")
    print("  1. 超时按「最慢那一路」设,不是按各路之和")
    print("  2. 先划出关键路,其余走尽力模式,别让可选项拖垮整条流程")
    print("  3. 缺了的分支要在下游提示词里显式标注,别让模型以为查过了")


if __name__ == "__main__":
    main()
四种收口情况的实跑结果
【三路全成功(严格模式)】
  ✅ order        order 的结果
  ✅ retrieve     retrieve 的结果
  ✅ websearch    websearch 的结果
  → 本次可以继续往下走
  耗时 0.55s —— 串行要 1.10s,并行只花最慢那一路的时间

【一路失败(严格模式)】
  ✅ order        order 的结果
  ✅ retrieve     retrieve 的结果
  ❌ websearch    websearch 调用失败
  → 本次不能继续往下走

【一路失败(尽力模式,关键路=retrieve+order)】
  ✅ order        order 的结果
  ✅ retrieve     retrieve 的结果
  ❌ websearch    websearch 调用失败
  → 本次可以继续往下走

【关键路失败(尽力模式)】
  ✅ retrieve     retrieve 的结果
  ✅ websearch    websearch 的结果
  ❌ order        order 调用失败
  ⛔ 关键路缺失:order
  → 本次不能继续往下走

收口的三个决定:
  1. 超时按「最慢那一路」设,不是按各路之和
  2. 先划出关键路,其余走尽力模式,别让可选项拖垮整条流程
  3. 缺了的分支要在下游提示词里显式标注,别让模型以为查过了

实跑里最值得注意的是第一段:三路耗时 0.20 + 0.35 + 0.55 秒,串行要 1.10 秒,并行只花了 0.55 秒——也就是最慢那一路的时间。这条规律直接决定了超时该怎么设:

超时按「最慢那一路」设,不是按各路之和 按各路之和设,等于给了一个根本用不上的宽松上限,某一路卡死时整条流程要等到天荒地老。按最慢那一路设(再留一点余量),卡死的那一路才会被及时切掉。
收口策略规则什么时候用
严格模式任一路失败,整体失败缺了就答不了,比如必须拿到订单状态才能回复
尽力模式关键路必须成功,其余缺了标注为空多一份资料更好、少一份也能答

看第三、四段实跑的差别:同样是一路失败,可选路挂了尽力模式能继续走;关键路挂了,尽力模式也得停。所以用尽力模式之前,必须先把关键路划出来——不划就等于全都是可选的,那才是真的危险。

3.3 循环:三道闸缺一不可

纯 DAG 不能表达「改到满意为止」,所以平台会额外给循环/迭代节点。代价是 DAG 的天然保证没了——它不再必然终止。

画布上一旦出现回边,DAG 的「必然终止」就没了 所以循环节点必须同时挂三道闸,少一道都可能跑出一条停不下来、或者停下来但白烧钱的流程 「改到满意为止」的形状 生成草稿 打分评审 达标,输出 不达标,重来 那条虚线就是回边。它让流程有了表达能力, 也让「一定会停」这个保证消失了。 三道闸,缺一不可 1 最大轮次 —— 保证一定会停 只有这一道时,流程能跑完,但可能转满全程一分没涨 2 预算上限 —— 保证停之前不会烧穿 每轮都在涨分也要拦:涨分不等于值这个价 3 无进展判定 —— 保证不做无用功 连续 N 轮增益低于阈值就停,别等轮次耗尽 实跑四种走向: 场景一 三轮收敛 0.62→0.78→0.86,闸未触发 场景二 原地踏步 第3轮被无进展闸拦下 场景三 一直在涨 0.900 撞上成本闸 0.800 场景四 只设轮次闸 转满 8 轮,花 2.00,0 涨幅
图⑤ 循环的三道闸:轮次、预算、无进展
循环三道闸:轮次、预算、无进展
"""循环与防护:让流程能回头,但不能永远回头。

纯 DAG 不能表达「改到满意为止」这类需求,所以平台会额外提供
循环/迭代节点。代价是 DAG 的天然保证没了 —— 它不再必然终止。

所以只要画布上出现回边,就必须同时给出三道闸:最大轮次、
预算上限、无进展判定。少一道都可能跑出一条停不下来的流程。

纯标准库,直接 python3 loop_guard.py 就能跑。
"""


class LoopBudget:
    """三道闸门合在一起:轮次、成本、进展。"""

    def __init__(self, max_rounds=5, max_cost=1.0, min_gain=0.02,
                 patience=2):
        self.max_rounds = max_rounds      # 闸一:最多转几圈
        self.max_cost = max_cost          # 闸二:累计花费上限
        self.min_gain = min_gain          # 闸三:单轮最小进步
        self.patience = patience          # 连续几轮没进步就停
        self.round = 0
        self.cost = 0.0
        self.stale = 0
        self.best = None

    def should_stop(self):
        if self.round >= self.max_rounds:
            return "达到最大轮次 %d" % self.max_rounds
        if self.cost >= self.max_cost:
            return "累计成本 %.3f 触达上限 %.3f" % (self.cost, self.max_cost)
        if self.stale >= self.patience:
            return "连续 %d 轮进步不足 %.3f" % (self.stale, self.min_gain)
        return None

    def record(self, score, cost):
        """登记一轮的结果,更新三道闸的状态。"""
        self.round += 1
        self.cost += cost
        if self.best is None:
            gain = score
        else:
            gain = score - self.best
        if self.best is None or score > self.best:
            self.best = score
        if gain < self.min_gain:
            self.stale += 1
        else:
            self.stale = 0
        return gain


def refine_loop(scores, costs, **kw):
    """模拟「生成→打分→不满意就重来」的循环。

    scores/costs 是预先准备好的每轮结果,保证输出可复现。
    """
    budget = LoopBudget(**kw)
    print("  轮次  本轮分  单轮增益  累计成本  状态")
    for score, cost in zip(scores, costs):
        stop = budget.should_stop()
        if stop:
            print("  —— 停止:%s" % stop)
            return budget
        gain = budget.record(score, cost)
        print("  %3d   %.3f   %+.3f     %.3f" %
              (budget.round, score, gain, budget.cost))
    stop = budget.should_stop()
    print("  —— 停止:%s" % (stop or "候选轮次用尽"))
    return budget


def main():
    print("【场景一】三轮就收敛,闸门没被触发")
    refine_loop([0.62, 0.78, 0.86, 0.87, 0.87],
                [0.12, 0.12, 0.12, 0.12, 0.12])

    print("\n【场景二】分数原地踏步,被「无进展」闸拦下")
    refine_loop([0.55, 0.56, 0.565, 0.57, 0.57],
                [0.12, 0.12, 0.12, 0.12, 0.12])

    print("\n【场景三】每轮都在涨,但先撞上成本闸")
    refine_loop([0.40, 0.55, 0.68, 0.79, 0.88],
                [0.30, 0.30, 0.30, 0.30, 0.30], max_cost=0.8)

    print("\n【场景四】只设了轮次闸,没设成本和进展闸")
    b = refine_loop([0.5] * 8, [0.25] * 8,
                    max_rounds=8, max_cost=999, min_gain=0.0)
    print("  结果:转满 %d 轮,花掉 %.2f,分数一点没涨" % (b.round, b.cost))
    print("  ⚠ 这就是「能跑完但白烧钱」的典型形态,靠轮次闸兜不住")

    print("\n只要画布上出现回边,三道闸缺一不可:")
    print("  闸一 最大轮次 —— 保证一定会停")
    print("  闸二 预算上限 —— 保证停之前不会烧穿")
    print("  闸三 无进展判定 —— 保证不做无用功")


if __name__ == "__main__":
    main()
四种走向的实跑结果
【场景一】三轮就收敛,闸门没被触发
  轮次  本轮分  单轮增益  累计成本  状态
    1   0.620   +0.620     0.120
    2   0.780   +0.160     0.240
    3   0.860   +0.080     0.360
    4   0.870   +0.010     0.480
    5   0.870   +0.000     0.600
  —— 停止:达到最大轮次 5

【场景二】分数原地踏步,被「无进展」闸拦下
  轮次  本轮分  单轮增益  累计成本  状态
    1   0.550   +0.550     0.120
    2   0.560   +0.010     0.240
    3   0.565   +0.005     0.360
  —— 停止:连续 2 轮进步不足 0.020

【场景三】每轮都在涨,但先撞上成本闸
  轮次  本轮分  单轮增益  累计成本  状态
    1   0.400   +0.400     0.300
    2   0.550   +0.150     0.600
    3   0.680   +0.130     0.900
  —— 停止:累计成本 0.900 触达上限 0.800

【场景四】只设了轮次闸,没设成本和进展闸
  轮次  本轮分  单轮增益  累计成本  状态
    1   0.500   +0.500     0.250
    2   0.500   +0.000     0.500
    3   0.500   +0.000     0.750
    4   0.500   +0.000     1.000
    5   0.500   +0.000     1.250
    6   0.500   +0.000     1.500
    7   0.500   +0.000     1.750
    8   0.500   +0.000     2.000
  —— 停止:达到最大轮次 8
  结果:转满 8 轮,花掉 2.00,分数一点没涨
  ⚠ 这就是「能跑完但白烧钱」的典型形态,靠轮次闸兜不住

只要画布上出现回边,三道闸缺一不可:
  闸一 最大轮次 —— 保证一定会停
  闸二 预算上限 —— 保证停之前不会烧穿
  闸三 无进展判定 —— 保证不做无用功

四个场景里,最值得看的是第四个:只设了轮次闸,结果转满 8 轮、花掉 2.00、分数一点没涨。这条流程「能跑完」,监控上也不会报警,但它每一次执行都在白烧钱。

⛔ 三道闸各管一件事 闸一 最大轮次——保证一定会停;闸二 预算上限——保证停之前不会烧穿(注意场景三:每轮都在涨分也照样拦,涨分不等于值这个价);闸三 无进展判定——保证不做无用功。只要画布上出现回边,三道一起挂。

04完整案例:四条真实形状的流程

每条都标出它的分层、关键路、以及最容易出事的那个节点

4.1 客服问答:分类 → 并行取数 → 组装 → 生成

这是最常见的形状,也是本页示例流程的原型。用户问「我买的杯子坏了,能退吗」,流程要同时做两件事:查知识库拿退换货政策,查订单系统拿这单的状态。

节点
1classify / retrievequestionintent / docs
2lookup_orderquestion, intentorder
3composedocs, order, intentprompt
4generatepromptanswer

关键路lookup_order。退换货问题里订单状态缺了就没法答——已签收和未发货的处理方式完全不同,猜错比不答更糟。可选路retrieve,政策文档缺了还能给个保守答复并转人工。

这条流程最容易出事的节点是 compose 因为它是唯一一个同时读三个上游变量的节点。上游任何一个改了输出名、改了类型、或者返回了 None,都在这里爆出来。排查编排故障时,先看汇聚节点。

4.2 批量内容生产:扇出到每一条数据

几千个商品,每个都要生成标题和卖点。形状上是一个循环包住一条短流程,但这里的循环是「遍历」而不是「改到满意为止」——轮次由数据条数决定,不由质量决定。

要点做法不这么做会怎样
并发度限制同时在跑的条数几千条一起发,立刻触发限流,全部失败
失败隔离单条失败不影响其他条一条挂掉整批回滚,前面几小时白跑
断点续跑记录已完成的条目 id中途中断只能从头再来
结果落盘时机每条跑完就写,不要攒到最后进程一挂,内存里的结果全丢

这个场景里没有关键路的概念——每条数据都是独立的一次执行,某条失败就单独重跑那一条。这也是它和客服问答最大的结构差别。

4.3 文档审阅:带回边的「改到满意为止」

生成草稿 → 打分评审 → 不达标就带着评语重新生成。这是回边的典型用法,也是三道闸必须全挂的场景。

设成什么拦的是哪种情况
最大轮次3–5 轮兜底,保证一定会停
预算上限按单次任务能接受的成本每轮都在涨分、但已经不值这个价
无进展判定连续 2 轮增益 < 阈值分数原地踏步,继续转纯属浪费

实跑的场景二就是被无进展闸拦下的:分数 0.550 → 0.560 → 0.565,每轮只涨零点几个百分点,第 3 轮就停了。如果只设轮次闸,它会一直转到第 5 轮,多花两轮的钱换来 0.005 的提升。

⚠️ 评审节点不要用同一个模型给自己打分 自己生成、自己打分,容易出现「越改越自信、分数越打越高」但实际质量没变的情况。要么换一个模型做评审,要么给评审节点一套明确的评分标准而不是让它凭感觉给分。

4.4 多路资料汇总:尽力模式的典型场景

写一份行业简报,要同时查内部知识库、公开网页、以及数据库里的历史数据。三路里网页搜索最不稳定,但它缺了简报照样能写。

这正好对应实跑的第三段:关键路 = retrieve + order,websearch 挂了照常继续。但有一个细节必须做到:

⛔ 缺了的分支要在下游提示词里显式标注 降级为 "(本次未取到外部资料)" 而不是空字符串,是为了让模型知道这一路是空的。如果只是给个空串,模型会以为你查过了、确实什么都没有,然后自信地写出「经查,该领域近期无重大进展」——这就从「资料缺失」变成了「编造结论」,性质完全不同。

4.5 四条流程横向看

流程形状有无回边最该防的事
客服问答分类 + 并行 + 汇聚汇聚节点拿到 None 或类型不对
批量生产遍历 + 短流程遍历,非质量回边并发打爆限流、中断后无法续跑
文档审阅生成 + 评审回边只设轮次闸,白转白花钱
资料汇总多路并行 + 尽力收口缺失分支没标注,模型编结论

05骨架模板:把画布变成能进 git 的东西

结构导出、指纹比对、变更清单——让流程的每次改动都能被看见

5.1 为什么要导出结构

平台的画布文件是私有格式,换平台带不走——这一点在上一页的锁定度模型里已经量化过,工作流编排结构的迁移难度是 0.7,仅次于分发渠道和会话历史。

但有个东西是通用的:编排的语义。有哪些节点、每个节点读什么写什么、连线怎么连——这层信息与平台无关,可以自己存一份。

先把期望说清楚:这不是「一键迁移」 没有一键迁移。导出结构的目的是:迁移时有一份准确的结构说明,而不是对着截图重画;以及出故障时,能确认「上线那天的流程到底长什么样」。

5.2 导出与指纹

把编排结构导成与平台无关的 JSON,并算指纹
"""把编排结构导出成自己能版本化的格式。

平台的画布文件是私有格式,换平台带不走。但编排的**语义**是通用的:
有哪些节点、每个节点读什么写什么、连线怎么连。把这层语义抽出来
存成自己的 JSON,进 git,就等于给流程上了保险。

这不是为了「以后一键迁移」——没有一键迁移。是为了迁移时
有一份准确的结构说明,而不是对着截图重画。

纯标准库,直接 python3 workflow_export.py 就能跑。
"""

import json
import hashlib

from dag_engine import build_demo


def export(wf):
    """把 Workflow 对象导成与平台无关的结构描述。"""
    nodes = []
    for name in wf.topo_order():
        n = wf.nodes[name]
        nodes.append({
            "id": name,
            "reads": list(n.reads),
            "writes": n.writes,
        })

    # 边由变量依赖推出来,不需要单独维护一份
    edges = []
    for n in nodes:
        for var in n["reads"]:
            src = wf.producer.get(var)
            if src:
                edges.append({"from": src, "to": n["id"], "via": var})

    doc = {
        "schema": "workflow-structure/v1",
        "nodes": nodes,
        "edges": edges,
        "layers": wf.parallel_layers(),
        "entry_vars": sorted({
            v for n in nodes for v in n["reads"]
            if v not in wf.producer
        }),
        "exit_vars": sorted({
            n["writes"] for n in nodes
            if not any(n["writes"] in m["reads"] for m in nodes)
        }),
    }
    return doc


def fingerprint(doc):
    """给结构算一个指纹,用来发现「画布被人改了但没人说」。"""
    payload = json.dumps(
        {"nodes": doc["nodes"], "edges": doc["edges"]},
        sort_keys=True, ensure_ascii=False).encode("utf-8")
    return hashlib.sha256(payload).hexdigest()[:16]


def diff(old, new):
    """比较两份导出,列出结构层面的变化。"""
    o = {n["id"]: n for n in old["nodes"]}
    n_ = {n["id"]: n for n in new["nodes"]}

    changes = []
    for nid in sorted(set(n_) - set(o)):
        changes.append("新增节点 %s" % nid)
    for nid in sorted(set(o) - set(n_)):
        changes.append("删除节点 %s" % nid)
    for nid in sorted(set(o) & set(n_)):
        if o[nid]["reads"] != n_[nid]["reads"]:
            changes.append("节点 %s 的输入变了:%s%s"
                           % (nid, o[nid]["reads"], n_[nid]["reads"]))
        if o[nid]["writes"] != n_[nid]["writes"]:
            changes.append("节点 %s 的输出变了:%s%s"
                           % (nid, o[nid]["writes"], n_[nid]["writes"]))
    return changes


def main():
    wf = build_demo()
    doc = export(wf)

    print("导出的结构(节选):")
    print(json.dumps(
        {k: doc[k] for k in ("schema", "entry_vars", "exit_vars", "layers")},
        ensure_ascii=False, indent=2))

    print("\n节点清单:")
    for n in doc["nodes"]:
        print("  %-14s reads=%-32s writes=%s"
              % (n["id"], ",".join(n["reads"]), n["writes"]))

    print("\n连线(由变量依赖自动推出,共 %d 条):" % len(doc["edges"]))
    for e in doc["edges"]:
        print("  %-14s --%s--> %s" % (e["from"], e["via"], e["to"]))

    fp1 = fingerprint(doc)
    print("\n结构指纹:%s" % fp1)

    # 模拟有人在画布上偷偷改了一个节点的输入
    wf2 = build_demo()
    wf2.nodes["compose"].reads = ["docs", "intent"]   # 不再读 order
    doc2 = export(wf2)
    fp2 = fingerprint(doc2)

    print("\n有人改了画布之后:")
    print("  新指纹:%s%s)" % (fp2, "一致" if fp1 == fp2 else "已变化"))
    for c in diff(doc, doc2):
        print("  • %s" % c)

    print("\n把这份 JSON 提交进 git,好处有三:")
    print("  1. 画布任何结构变动都能在代码评审里被看见")
    print("  2. 迁移时照着它重建,而不是对着截图猜")
    print("  3. 出故障时能确认「上线那天的流程到底长什么样」")


if __name__ == "__main__":
    main()
导出结果、连线推导、以及改动前后的指纹比对
导出的结构(节选):
{
  "schema": "workflow-structure/v1",
  "entry_vars": [
    "question"
  ],
  "exit_vars": [
    "answer"
  ],
  "layers": [
    [
      "classify",
      "retrieve"
    ],
    [
      "lookup_order"
    ],
    [
      "compose"
    ],
    [
      "generate"
    ]
  ]
}

节点清单:
  classify       reads=question                         writes=intent
  retrieve       reads=question                         writes=docs
  lookup_order   reads=question,intent                  writes=order
  compose        reads=docs,order,intent                writes=prompt
  generate       reads=prompt                           writes=answer

连线(由变量依赖自动推出,共 5 条):
  classify       --intent--> lookup_order
  retrieve       --docs--> compose
  lookup_order   --order--> compose
  classify       --intent--> compose
  compose        --prompt--> generate

结构指纹:8be845a9d6cc2e97

有人改了画布之后:
  新指纹:cf4c557af96ed3d0(已变化)
  • 节点 compose 的输入变了:['docs', 'order', 'intent'] → ['docs', 'intent']

把这份 JSON 提交进 git,好处有三:
  1. 画布任何结构变动都能在代码评审里被看见
  2. 迁移时照着它重建,而不是对着截图猜
  3. 出故障时能确认「上线那天的流程到底长什么样」

导出的结构里有几个字段值得单独说:

字段内容用途
nodes每个节点的 reads / writes迁移时照着它在新平台上重建
edges由变量依赖自动推出不用单独维护,也就不会和实际不一致
layers并行分层结果一眼看出哪些节点本可并行却被串起来了
entry_vars没有任何节点写过的变量这就是流程的输入契约
exit_vars写了之后没人读的变量正常情况下只该有一个(最终输出);多出来的就是废节点

5.3 指纹:发现「画布被人改了但没人说」

把节点和边序列化后算一个哈希,就得到结构指纹。实跑里模拟了一次偷偷修改——把 compose 的输入里的 order 去掉:

时间点指纹差异
改动前8be845a9d6cc2e97
改动后cf4c557af96ed3d0节点 compose 的输入变了:['docs','order','intent']['docs','intent']

注意这个改动的性质:流程照样能跑,不报任何错,只是从此以后所有回复都不再带订单信息了。没有指纹比对,这种变化要等到用户投诉才会被发现。

5.4 落地节奏

把这套东西接进日常工作,只需要三步:

1每次上线前导出

把结构 JSON 和指纹一起提交进 git,作为这次上线的记录。

2评审时看 diff

结构变更会以清晰的文字出现在代码评审里,而不是埋在一张截图里。

3定期比对线上

把线上实际结构导出来和 git 里的比指纹,不一致就说明有人直接改了线上画布。

4故障时回放

对着故障时间点那一版结构排查,不用猜当时流程是什么样。

⛔ 一个容易被忽略的收益 exit_vars 里多出来的变量,等于自动帮你找出了写了没人读的废节点。这些节点每次执行都在花钱,却对结果毫无贡献。定期扫一遍,通常都能清掉几个。

06易错点:九个让流程静默变差的做法

最危险的不是报错的那些,是不报错但答案悄悄变差的那些

① 以为画布上的位置决定执行顺序

错在哪:把节点从左往右排好,就以为它们会按这个顺序跑。实际顺序由依赖关系算出来,跟位置无关。

翻车的样子:调了半天位置,顺序还是不对;或者两个本可并行的节点被无意中串成了一条线,白白多花一倍时间。

怎么防:想让 B 等 A,就让 B 读 A 写的变量。想让两个节点并行,就确认它们没有互相读写。

② 给所有节点统一开「重试 3 次」

错在哪bad_requestauth_failedcontent_filtered 这类错误,重试一百次结果相同。

翻车的样子:一次失败被放大成四次,延迟翻四倍,配额也白烧四份。

怎么防:按错误类型分流。只重试 timeout / rate_limited / upstream_5xx 这三类。

③ 退避不加抖动

错在哪:所有被同一次故障打挂的请求,会在同一时刻一起重试。

翻车的样子:下游刚缓过来就被第二波齐射打死,小故障升级成大事故。

怎么防:退避时间上加随机抖动,把重试摊开到一段区间里。

④ 并行超时按「各路之和」设

错在哪:并行的总耗时等于最慢那一路,不是各路相加。实跑数据:三路 0.20 + 0.35 + 0.55,并行只花 0.55 秒。

翻车的样子:超时设成 1.10 秒,某一路卡死时白等到超时上限才切断。

怎么防:按最慢那一路的正常耗时设,再留一点余量。

⑤ 用尽力模式但没划关键路

错在哪:没划关键路,等于所有分支都是可选的。

翻车的样子:订单查询挂了,流程照样往下走,模型在没有订单信息的情况下编了一个退货结论。

怎么防:用尽力模式之前先回答一个问题——哪一路缺了就不能答?答不上来就说明该用严格模式。

⑥ 降级成空字符串

错在哪:把失败分支降级为 "",模型会以为你查过了、确实什么都没有。

翻车的样子:从「资料缺失」变成「编造结论」——模型自信地写出「经查,该领域近期无重大进展」。这两者性质完全不同。

怎么防:降级值必须显式说明状态,比如 "(本次未取到外部资料)",让下游知道这一路是空的。

⑦ 循环只设最大轮次

错在哪:轮次闸只保证「一定会停」,不保证「停之前有价值」。

翻车的样子:实跑场景四——转满 8 轮、花掉 2.00、分数一点没涨。监控上不会报警,因为它确实跑完了。

怎么防:三道闸一起挂。预算上限拦烧穿,无进展判定拦白转。

⑧ 两个节点写同一个变量

错在哪:变量表是全局共享且可覆盖的,后写的静默盖掉先写的。

翻车的样子:实跑里 lookup_orderlookup_ticket 都写 result,下游拿到的是工单信息,但看上去像模型答错了。不报错,最难查。

怎么防:变量名带节点前缀(retrieve_docs),并在引擎或评审环节禁止重复写入。

⑨ 直接在线上画布改流程

错在哪:改动不经过评审,也没有记录。参考实跑里那次修改——把 compose 的一个输入去掉,流程照样跑、不报任何错,只是从此回复里不再带订单信息。

翻车的样子:等用户投诉才发现,而且没人说得清是哪天改的、改了什么。

怎么防:导出结构进 git,比指纹。定期把线上实际结构导出来和 git 里的比一次。

⛔ 九条压成三句 顺序由依赖决定,不由位置决定不报错的故障比报错的危险(类型漂移、同名覆盖、空降级、白转循环,全都不报错);画布也是代码,要进 git、要评审、要能回放

07自测题

点击题目展开答案;能把这 16 题说清楚,这一页就通了

一、节点与连线
一个编排节点由哪三段构成?不同类型的节点差别在哪一段?

reads(声明读哪些变量)、fn(真正干活的逻辑)、writes(声明写进哪个变量)。大模型节点、代码节点、知识库节点、HTTP 节点的差别只在 fn 里干了什么,结构完全一样。

画布上那根连线,和真正的依赖关系是什么关系?

连线只是依赖关系的可视化,不是依赖关系本身。真正的依赖来自变量:compose 读了 docs,而 docsretrieve 写的,那条边就已经存在了。所以引擎可以从 reads/writes 自动推出全部连线。

「连线漏了」的本质是什么?怎么排查最快?

本质是下游读了一个没人写过的变量名。所以排查时别盯着画布找哪根线没连,直接查变量名对不对得上,快得多。报错信息里那三种可能:上游改了输出名、连线漏了、或这个节点被排到了上游前面。

为什么强调「一个节点只写一个变量」?

因为变量表是全局共享、可被覆盖、没有类型检查的,只有一个命名空间。一个节点写多处,等于在全局里多埋几颗雷,撞名就静默覆盖。

二、执行顺序与并行
执行顺序是由什么决定的?拓扑排序顺带解决了什么问题?

依赖关系算出来的,跟你在画布上画的位置无关。拓扑排序统计入度、逐个消解,顺带解决环检测:排完还有节点没排进去,说明它们入度永远减不到 0,也就是存在回边。

示例流程里 classify 和 retrieve 为什么能并行?

因为两个都只读 question,谁也不读谁的输出,互不依赖,所以同属第 1 层。注意 retrieve 在画布上可能被画得很靠后,但它第一层就能开跑。

怎样会不小心把两个本可并行的节点变成串行?

给后面那个节点加一个它根本不需要的输入。一旦它读了前一个节点写的变量,依赖就产生了,并行就没了。所以 reads 要写得精确,别图省事把所有上游变量都列上。

并行的总耗时怎么算?超时该按什么设?

总耗时等于最慢那一路。实跑数据:三路 0.20+0.35+0.55 秒,串行要 1.10 秒,并行只花 0.55 秒。超时应按最慢那一路的正常耗时再留余量来设,按各路之和设等于给了个用不上的宽松上限,某路卡死时白等。

三、变量表上的故障
变量表的哪三个特点造成了编排里最难查的 bug?

全局共享(只有一个命名空间,撞名就覆盖)、可被覆盖(写入即赋值,不提示重名)、没有类型检查(今天是 list 明天是 str 也不管)。第四个特点「只增不改」则是好事,让排障时能完整回放中间值。

三类故障分别是什么?哪两类不报错?

类型漂移(docs 从 list 变 str,len() 从 2 变 7,不报错)、改名断链(上游改叫 chunks,下游还读 docs,会报错但指向不明)、同名覆盖(订单信息被工单盖掉,不报错)。不报错的那两类最危险——答案只是悄悄变差。

防住这三类故障的三条防线是什么?

① 变量名带节点前缀,如 retrieve_docs,天然不会撞名;② 每个节点只写一个变量;③ 上线前跑一次作用域体检,清掉「写了没人读」的节点——那种节点要么是废的,要么是下游连线漏了。

四、失败、循环与版本
哪些错误该重试、哪些不该?为什么不能统一开重试 3 次?

该重试:timeoutrate_limitedupstream_5xx——再试一次可能就好了。不该重试:bad_requestauth_failedcontent_filtered——参数或凭证本身就错,试一百次结果相同。统一开重试会把一次失败放大成四次,延迟翻四倍,配额白烧四份。

退避为什么必须加抖动?

不加抖动时,被同一次故障打挂的所有请求会在同一时刻一起重试,把刚缓过来的下游再次打死,这叫重试风暴。对比:不带抖动 [0.4, 0.8, 1.6, 3.2, 6.4],带抖动 [0.38, 0.75, 1.29, 3.2, 6.13]——抖动把重试摊开到一段区间里。

严格模式和尽力模式怎么选?用尽力模式前必须先做什么?

缺了就答不了的用严格模式(任一路失败整体失败);多一份更好、少一份也能答的用尽力模式。用尽力模式前必须先划出关键路——不划等于所有分支都可选,关键路挂了流程照样往下走,模型就会在缺数据的情况下编结论。

失败分支为什么不能降级成空字符串?

因为模型会以为你查过了、确实什么都没有,然后自信地写出「经查,该领域近期无重大进展」。这就从「资料缺失」变成了「编造结论」。降级值要显式说明状态,比如「(本次未取到外部资料)」。

循环节点为什么必须挂三道闸?只设轮次闸会怎样?

因为回边一出现,DAG「必然终止」的保证就没了。三道闸:最大轮次保证一定会停,预算上限保证停之前不烧穿,无进展判定保证不做无用功。只设轮次闸的后果见实跑场景四:转满 8 轮、花掉 2.00、分数一点没涨,而且监控不会报警,因为它确实跑完了。

结构指纹能发现什么?举一个「不报错但很严重」的例子。

能发现画布被人改了但没人说。实跑例子:把 compose 的输入里的 order 去掉,指纹从 8be845a9d6cc2e97 变成 cf4c557af96ed3d0。这个改动流程照样跑、不报任何错,只是从此回复不再带订单信息——没有指纹比对,要等用户投诉才会发现。

术语表

这一页出现的概念,按排查故障时被用到的频率排

术语英文 / 写法说明
节点node编排的基本单位,由 reads / fn / writes 三段构成。不同类型的节点差别只在 fn。
变量表state所有节点共用的一份字典。全局共享、可被覆盖、无类型检查、只增不改——前三条是故障来源,最后一条让排障可回放。
连线edge依赖关系的可视化。真正的边由「谁读了谁写的变量」自动推出,不需要手工维护
拓扑排序topological sort按入度消解算出执行顺序的算法,顺带完成环检测
并行分层parallel layers把节点按「依赖是否已满足」分层,同层互不依赖可并行。能否并行由依赖决定,不由画布位置决定。
回边back edge让流程能回头的那条线。它带来「改到满意为止」的表达能力,代价是不再保证必然终止
类型漂移type drift同一个变量在不同分支里是不同类型。不报错,只是答案悄悄变差,最难查。
同名覆盖variable shadowing两个节点写同一个变量,后写的静默盖掉先写的。不报错,看上去像模型答错了。
改名断链broken reference上游改了输出名,下游还读旧名字。会报错,但报错指向不明。
指数退避exponential backoff每次重试等待时间翻倍并设上限,避免密集重试压垮下游。
抖动jitter在退避时间上加随机量。不加抖动会造成重试风暴——所有请求同一刻齐射,把刚缓过来的下游再打死。
可重试错误retryabletimeoutrate_limitedupstream_5xx——再试一次可能就好了。
不可重试错误non-retryablebad_requestauth_failedcontent_filtered——重试只会把一次失败放大成四次。
关键路required branch缺了就不能继续的分支。用尽力模式前必须先划出来,不划等于所有分支都可选。
严格模式 / 尽力模式all-or-nothing / best-effort前者任一路失败即整体失败;后者关键路必须成功、其余缺了标注为空。
降级fallback主路径失败后给出的保守答案。必须显式说明状态,降级成空字符串会让模型把「资料缺失」当成「确实没有」。
三道闸loop budget循环节点的防护:最大轮次(一定会停)、预算上限(不烧穿)、无进展判定(不做无用功)。缺一不可。
结构指纹structure fingerprint节点与边序列化后的哈希。用来发现画布被人改了但没人说,尤其是那种不报错的改动。
输入契约 / 出口变量entry_vars / exit_vars前者是没有任何节点写过的变量,也就是流程的输入;后者写了没人读,正常只该有一个,多出来的就是废节点
三句话收口 顺序由依赖决定,不由位置决定不报错的故障比报错的危险画布也是代码,要进 git、要评审、要能回放