乐于分享
好东西不私藏

S1-源码学习篇-07-工作流引擎:让 AI 按"流程图"干活——万悟 DAG 编排拆解

S1-源码学习篇-07-工作流引擎:让 AI 按"流程图"干活——万悟 DAG 编排拆解

工作流引擎:让 AI 按"流程图"干活——万悟 DAG 编排源码拆解

📌 本文是《从零吃透企业级 AI 平台:元景万悟源码学习手记》系列第一季·源码学习篇的第 7 篇🎯 读完本文你将:① 理解 DAG 工作流引擎的核心机制 ② 掌握条件分支、并行执行、错误处理的实现原理 ③ 能在万悟中设计并调试一个多节点工作流⏱️ 预计阅读时间:35 分钟 | 动手实践:80 分钟💻 前置要求:已完成第 5、6 篇,理解 Agent 与 MCP 工具调用机制


一、这篇文章要解决什么问题?

第 5 篇的 Agent 很强大,但它有一个根本问题:不确定性

Agent 的推理链是 LLM 动态决定的——同样的问题问两次,它可能走不同的路径、调不同的工具、甚至得出不同的结论。这在"闲聊问答"场景没问题,但在企业级流程中是致命的:

"每周一早上 9 点,自动拉取上周销售数据 → 生成报表 → 如果低于目标就发预警邮件 → 否则发周报 → 归档到知识库"

这个流程:

  • 步骤固定:不需要 LLM "思考"下一步做什么
  • 顺序严格:必须先拉数据,再生成报表,不能乱
  • 分支明确:低于目标走 A 路径,否则走 B 路径
  • 必须可靠:不能因为 LLM "心情不好"就跳过某一步

工作流(Workflow)= 预定义的 DAG(有向无环图)+ 确定性执行引擎

Agent:  ”你看着办”(LLM 动态决策,灵活但不确定)工作流:  ”按这个流程图走”(预定义路径,确定但不灵活)

💡 经验法则

  • 步骤固定、分支明确 → 用工作流
  • 步骤不确定、需要推理 → 用 Agent
  • 固定流程中某一步需要推理 → 工作流 + 内嵌 Agent 节点

万悟的工作流引擎支持可视化拖拽编排,底层是一个 DAG 执行器。今天我们把它拆开看。


二、核心概念:用大白话讲清楚

2.1 DAG:有向无环图

💡 类比:工作流就像做菜的流程图

  • 洗菜 → 切菜 → 炒菜(串行)
  • 同时:煮饭(并行)
  • 如果菜是辣的 → 多放糖;否则 → 正常调味(条件分支)
  • 所有步骤完成 → 上桌(汇聚)

关键约束:不能成环。你不能"炒完菜再回去洗菜"。这就是"无环"的含义。

┌──→ [切菜] ──→ [炒菜] ──┐[洗菜] ─┤                        ├──→ [上桌]        └──→ [煮饭] ────────────┘                    │                    ▼              [如果辣?]              ╱        ╲           是            否           ▼              ▼       [多放糖]      [正常调味]

2.2 工作流的核心元素

元素
说明
万悟中的实现
节点(Node)
一个执行单元
LLM 调用、工具调用、代码执行、条件判断、子流程
边(Edge)
节点间的连接
数据流向 + 执行顺序
条件(Condition)
分支判断
基于上游节点输出的表达式
变量(Variable)
节点间传递的数据
{{node_1.output}}
 模板引用
触发器(Trigger)
工作流启动条件
手动、定时(Cron)、Webhook、事件

2.3 万悟支持的节点类型

节点类型
功能
典型用途
开始节点
接收输入参数
定义工作流入参 Schema
LLM 节点
调用大模型
文本生成、分类、摘要
工具节点
调用 MCP/内置工具
查数据库、发消息、调 API
代码节点
执行 Python/JS 代码
数据转换、格式处理
条件节点
if/else 分支
根据上一步结果走不同路径
循环节点
对列表逐项执行
批量处理
并行节点
多路同时执行
同时查多个数据源
子流程节点
嵌套另一个工作流
复用、模块化
结束节点
输出最终结果
定义工作流出参

🔒 代码节点安全——wga-sandbox 容器隔离

代码节点允许执行用户提供的 Python/JS 代码,这是一个高风险操作。万悟的做法是——不在主进程中直接执行,而是通过独立的 wga-sandbox 容器(:4096) 做沙箱隔离:

CODE_RUNNER_TYPE: sandbox

沙箱执行模式

CODE_RUNNER_ALLOW_RUN: "python"

仅允许 Python(禁用 JS)

CODE_RUNNER_TIMEOUT_SECONDS: 30

30 秒超时

CODE_RUNNER_MEMORY_LIMIT_MB: 512

512MB 内存上限

代码跑在隔离容器里,崩了只毁容器不拖垮 workflow-wanwu。同时通过 `WANWU_WGA_SANDBOX_IMAGE` 指定沙箱镜像版本,可独立升级。

2.4 工作流 vs Agent:什么时候用哪个?

维度
工作流
Agent
执行路径
预定义、确定性
LLM 动态决策
可预测性
✅ 高(每次走同一路径)
❌ 低(可能走不同路径)
灵活性
❌ 低(改流程要改图)
✅ 高(LLM 自适应)
调试难度
✅ 低(每步可断点)
❌ 高(推理链不确定)
适用场景
审批流、报表生成、ETL
开放问答、复杂推理
成本可控性
✅ 高(LLM 调用次数固定)
❌ 低(循环次数不确定)

📌 最佳实践:企业级应用中,80% 的场景应该用工作流,20% 用 Agent。工作流中需要"智能"的节点,嵌入一个 Agent 或 LLM 节点即可。


三、万悟是怎么实现的?(源码篇)

3.1 定位源码

⚠️ 诚实说明:万悟的工作流引擎代码不在本仓库——它在独立仓库 UnicomAI/wanwu-workflow(v0.2.0+)。本仓库(UnicomAI/wanwu)仅通过 docker-compose 拉起 workflow 容器(8998/8999),并在 bff-service 暴露工作流 API 与可视化画布。下面 3.2–3.6 的代码是教学重建版——用最小 Go 代码讲清 DAG 引擎的核心逻辑。真实实现请看 wanwu-workflow 仓库。

# 万悟工作流真实架构# 1. 独立仓库:UnicomAI/wanwu-workflow(DAG 引擎、节点执行器、调度器)# 2. 本仓库中的集成点:internal/bff-service/# BFF 网关(暴露工作流 API + 可视化画布)docker-compose*.yaml# 拉起 workflow 容器(8998/8999)# 工作流容器通过 gRPC 与 bff-service 通信# bff-service 负责工作流 CRUD、执行触发、状态查询# workflow 容器负责 DAG 解析、节点调度、执行

3.2 DAG 执行引擎

// 教学重建版(真实实现:UnicomAI/wanwu-workflow 独立仓库)// Execute 是工作流执行的主入口// 它按照 DAG 拓扑顺序执行所有节点func(e *Executor) Execute(ctx context.Context, wf *Workflow, input map[string]interface{}) (*ExecutionResult, error) {    // Step 1: 构建执行上下文(存储所有节点的输入/输出)    execCtx := NewExecutionContext(input)    // Step 2: 拓扑排序,确定执行顺序    // 将 DAG 转为层级结构:同一层的节点可以并行执行    layers, err := e.scheduler.TopologicalSort(wf.Nodes, wf.Edges)    if err != nil {        return nil, fmt.Errorf(”topological sort failed (cycle detected?): %w”, err)    }    // Step 3: 逐层执行    for layerIdx, layer := range layers {        log.Infof(”executing layer %d, %d nodes”, layerIdx, len(layer))        // 同一层的节点并行执行        var wg sync.WaitGroup        errCh := make(chan errorlen(layer))        for _, nodeID := range layer {            node := wf.GetNode(nodeID)            // 检查条件:如果上游条件不满足,跳过此节点            if !e.shouldExecute(node, execCtx) {                execCtx.SetNodeStatus(nodeID, NodeSkipped)                continue            }            wg.Add(1)            go func(n *Node) {                defer wg.Done()                // 执行节点(带重试)                output, err := e.executeNodeWithRetry(ctx, n, execCtx)                if err != nil {                    errCh <- fmt.Errorf(”node %s failed: %w”, n.ID, err)                    execCtx.SetNodeStatus(n.ID, NodeFailed)                    return                }                // 将输出存入上下文,供下游节点引用                execCtx.SetNodeOutput(n.ID, output)                execCtx.SetNodeStatus(n.ID, NodeSuccess)            }(node)        }        wg.Wait()        close(errCh)        // 检查是否有节点失败        for err := range errCh {            if wf.ErrorStrategy == ”fail_fast” {                return nil, err // 快速失败            }            // ”continue” 策略:记录错误但继续执行其他分支            log.Warnf(”node error (continuing): %v”, err)        }    }    // Step 4: 收集最终输出    return execCtx.GetFinalOutput(wf.EndNodeID), nil}

📌 关键洞察:整个引擎的核心就是拓扑排序 + 逐层并行。同一层(无依赖关系)的节点用 goroutine 并发执行,不同层(有依赖)严格串行。这既保证了正确性,又最大化了并行度。

3.3 拓扑排序与调度

// 教学重建版(真实实现:UnicomAI/wanwu-workflow 独立仓库)// TopologicalSort 将 DAG 节点分层// 返回 [][]string,每个内层 slice 是可以并行执行的节点 IDfunc(s *Scheduler) TopologicalSort(nodes []*Node, edges []*Edge) ([][]stringerror) {    // 构建邻接表和入度表    inDegree := make(map[string]int)    adjacency := make(map[string][]string)    for _, node := range nodes {        inDegree[node.ID] = 0    }    for _, edge := range edges {        adjacency[edge.From] = append(adjacency[edge.From], edge.To)        inDegree[edge.To]++    }    // BFS 分层(Kahn 算法)    var layers [][]string    queue := []string{}    // 入度为 0 的节点是第一层    for id, deg := range inDegree {        if deg == 0 {            queue = append(queue, id)        }    }    for len(queue) > 0 {        // 当前层的所有节点        layers = append(layers, queue)        var nextQueue []string        for _, nodeID := range queue {            for _, neighbor := range adjacency[nodeID] {                inDegree[neighbor]--                if inDegree[neighbor] == 0 {                    nextQueue = append(nextQueue, neighbor)                }            }        }        queue = nextQueue    }    // 检测环:如果还有节点入度 > 0,说明有环    for _, deg := range inDegree {        if deg > 0 {            return nil, fmt.Errorf(”cycle detected in workflow DAG”)        }    }    return layers, nil}

💡 算法复习:这就是经典的 Kahn 算法(BFS 拓扑排序)。时间复杂度 O(V+E)。万悟在保存工作流时就做一次环检测,拒绝含环的定义;执行时再排一次确定层级。

3.4 条件节点

// 教学重建版(真实实现:UnicomAI/wanwu-workflow 独立仓库)// Execute 评估条件表达式,决定走哪个分支func(c *ConditionNode) Execute(ctx *ExecutionContext) (*NodeOutput, error) {    // 条件表达式示例:{{node_2.output.total_sales}} < 100000    expression := c.Config.Expression    // 模板变量替换:将 {{node_2.output.total_sales}} 替换为实际值    resolved := ctx.ResolveTemplate(expression)    // 替换后:85000 < 100000    // 安全评估表达式(不用 eval,用白名单解析器)    result, err := c.evaluator.Evaluate(resolved)    if err != nil {        return nil, fmt.Errorf(”condition evaluation failed: %w”, err)    }    // 返回分支选择    if result.Bool() {        return &NodeOutput{Branch: ”true”}, nil   // 走”是”分支    }    return &NodeOutput{Branch: ”false”}, nil      // 走”否”分支}

📌 安全设计:注意万悟不用 eval() 执行条件表达式,而是用白名单解析器(只支持比较、逻辑运算、数学运算)。这防止了用户通过条件表达式注入恶意代码。

3.5 变量传递与模板解析

// 教学重建版(真实实现:UnicomAI/wanwu-workflow 独立仓库)// ResolveTemplate 解析模板变量引用// 输入:”销售额为 {{node_2.output.total}},目标为 {{workflow.input.target}}”// 输出:”销售额为 85000,目标为 100000”func(ctx *ExecutionContext) ResolveTemplate(template stringstring {    // 正则匹配 {{...}} 占位符    re := regexp.MustCompile(`\{\{(.+?)\}\}`)    return re.ReplaceAllStringFunc(template, func(match string) string {        // 去掉 {{ }}        path := strings.Trim(match, ”{} ”)        // 路径解析:node_2.output.total → 节点2的输出中的total字段        value, err := ctx.ResolvePath(path)        if err != nil {            return match // 解析失败,保留原文        }        return fmt.Sprintf(”%v”, value)    })}// ResolvePath 按路径查找变量值func(ctx *ExecutionContext) ResolvePath(path string) (interface{}, error) {    parts := strings.Split(path, ”.”)    switch parts[0] {    case ”workflow”:        // workflow.input.xxx → 工作流全局输入        return getNestedValue(ctx.globalInput, parts[2:])    case ”node”:        // node_2.output.xxx → 某个节点的输出        nodeID := parts[0] + ”.” + parts[1// ”node_2”        output, ok := ctx.nodeOutputs[nodeID]        if !ok {            return nil, fmt.Errorf(”node %s not executed yet”, nodeID)        }        return getNestedValue(output, parts[3:])    default:        return nil, fmt.Errorf(”unknown variable scope: %s”, parts[0])    }}

💡 设计亮点{{node_2.output.total}} 这种模板语法让用户在可视化编辑器中可以拖拽连线来传递数据,而不需要写代码。万悟前端在连线时自动生成这个模板字符串。

3.6 重试与错误处理

// 教学重建版(真实实现:UnicomAI/wanwu-workflow 独立仓库)// executeNodeWithRetry 带重试的节点执行func(e *Executor) executeNodeWithRetry(ctx context.Context, node *Node, execCtx *ExecutionContext) (*NodeOutput, error) {    maxRetries := node.RetryConfig.MaxRetries     // 默认 3    backoff := node.RetryConfig.InitialBackoff    // 默认 1s    var lastErr error    for attempt := 0; attempt <= maxRetries; attempt++ {        if attempt > 0 {            log.Infof(”retrying node %s, attempt %d/%d”, node.ID, attempt, maxRetries)            time.Sleep(backoff)            backoff *= 2 // 指数退避:1s → 2s → 4s        }        output, err := e.executeNode(ctx, node, execCtx)        if err == nil {            return output, nil        }        lastErr = err        // 判断是否可重试(网络超时可重试,参数错误不可重试)        if !isRetryable(err) {            break        }    }    return nil, fmt.Errorf(”node %s failed after %d retries: %w”, node.ID, maxRetries, lastErr)}

四、动手跑通(实践篇)

4.1 设计一个实际工作流

🎯 目标:构建"智能周报生成"工作流

[开始] → [查数据库:本周销售数据] → [LLM:生成分析摘要]                                         │                                    [条件:是否达标?]                                    ╱              ╲                                 是                 否                                 ▼                  ▼                          [LLM:写周报]      [LLM:写预警报告]                                 │                  │                                 ▼                  ▼                          [发邮件:周报]     [发邮件:预警]                                 ╲                ╱                                  ╲              ╱                                   ▼            ▼                              [归档到知识库]                                   │                                   ▼                                [结束]

4.2 在万悟中搭建

  1. 「工作流」→「新建」→ 拖入节点:   - 开始节点:定义输入参数 week_startweek_end   - 工具节点(MySQL):SELECT SUM(amount) FROM sales WHERE date BETWEEN '{{workflow.input.week_start}}' AND '{{workflow.input.week_end}}'   - LLM 节点:Prompt = 分析以下销售数据,给出关键洞察:{{node_1.output}}   - 条件节点:{{node_2.output.total}} >= {{workflow.input.target}}   - 两个 LLM 分支节点   - 工具节点(邮件)   - 工具节点(知识库归档)   - 结束节点
  2. 连线(拖拽)
  3. 配置每个节点的参数

▲ 万悟工作流可视化编辑器:拖拽节点、连线、配置参数

4.3 调试执行

  1. 点击「测试运行」
  2. 输入测试参数:week_start=2026-07-20week_end=2026-07-26target=100000
  3. 观察执行过程:
节点
状态
耗时
输出
开始
0ms
{week_start: "2026-07-20", ...}
查数据库
120ms
{total: 85000, orders: 342}
LLM 分析
2.3s
"本周销售额85000,环比下降12%..."
条件判断
1ms
false
(85000 < 100000)
写预警报告
1.8s
"⚠️ 本周未达标..."
发邮件
350ms
sent_to: boss@company.com
归档
200ms
doc_id: KB-2026-0727
结束
0ms
最终输出

📌 调试技巧:万悟支持单步执行——点击某个节点可以单独运行,查看它的输入/输出。这对定位"哪个节点出了问题"非常有用。

4.4 配置触发器

触发方式
配置
场景
手动
默认
测试、临时执行
定时(Cron)
0 9 * * 1
(每周一 9:00)
周报、日报
Webhook
生成 URL,外部系统 POST 触发
订单完成后触发
事件
监听平台内部事件
知识库更新后触发

4.5 对比实验

实验
操作
观察点
串行 vs 并行
将"查数据库"和"查CRM"改为并行节点
总耗时变化
有重试 vs 无重试
模拟工具节点超时
是否自动恢复
工作流 vs Agent
同一任务分别用工作流和 Agent 实现
确定性、耗时、成本

五、自己造一个 Mini 版(Deep Dive)

🎯 目标:用 Python 实现一个 80 行的 DAG 工作流引擎

# mini_workflow.py# 零依赖!仅用标准库实现 DAG 工作流引擎import jsonfrom collections import defaultdict, dequefrom concurrent.futures import ThreadPoolExecutor# ===== 1. 定义节点 =====class Node:    def __init__(self, id, func, deps=None):        self.id = id        self.func = func# 可调用对象        self.deps = deps or []# 依赖的节点 ID 列表# ===== 2. DAG 引擎 =====class WorkflowEngine:    def __init__(self):        self.nodes = {}        self.outputs = {}# node_id → output    def add_node(self, node):        self.nodes[node.id] = node    def _topo_layers(self):        ”””Kahn 算法分层”””        in_deg = {nid: 0 for nid in self.nodes}        adj = defaultdict(list)        for nid, node in self.nodes.items():            for dep in node.deps:                adj[dep].append(nid)                in_deg[nid] += 1        layers, queue = [], deque([n for n, d in in_deg.items() if d == 0])        while queue:            layer = list(queue)            layers.append(layer)            next_q = deque()            for nid in layer:                for nb in adj[nid]:                    in_deg[nb] -= 1                    if in_deg[nb] == 0:                        next_q.append(nb)            queue = next_q        if sum(in_deg.values()) > 0:            raise ValueError(”Cycle detected!”)        return layers    def run(self, global_input):        ”””执行工作流”””        self.outputs = {”__input__”: global_input}        layers = self._topo_layers()        for layer_idx, layer in enumerate(layers):            print(f”--- Layer {layer_idx}: {layer} ---”)# 同层并行执行            with ThreadPoolExecutor(max_workers=4as pool:                futures = {}                for nid in layer:                    node = self.nodes[nid]# 收集依赖节点的输出作为输入                    node_input = {dep: self.outputs[dep] for dep in node.deps}                    node_input[”__input__”] = global_input                    futures[nid] = pool.submit(node.func, node_input)                for nid, future in futures.items():                    self.outputs[nid] = future.result()                    print(f”  ✅ {nid} → {self.outputs[nid]}”)        return self.outputs# ===== 3. 定义工作流 =====engine = WorkflowEngine()# 节点函数def fetch_sales(ctx):    return {”total”: 85000, ”orders”: 342}def analyze(ctx):    data = ctx[”fetch_sales”]    return f”销售额{data['total']},订单{data['orders']}单,环比下降12%”def check_target(ctx):    data = ctx[”fetch_sales”]    return ”pass” if data[”total”] >= 100000 else ”fail”def write_report(ctx):    analysis = ctx[”analyze”]    return f”📊 周报:{analysis}。整体平稳。”def write_alert(ctx):    analysis = ctx[”analyze”]    return f”⚠️ 预警:{analysis}。建议加大推广力度。”def send_email(ctx):# 根据上游分支选择内容    content = ctx.get(”write_report”) or ctx.get(”write_alert”)    return f”邮件已发送: {content[:30]}...”# 注册节点(deps 定义依赖关系)engine.add_node(Node(”fetch_sales”, fetch_sales))engine.add_node(Node(”analyze”, analyze, deps=[”fetch_sales”]))engine.add_node(Node(”check_target”, check_target, deps=[”fetch_sales”]))engine.add_node(Node(”write_report”, write_report, deps=[”analyze”, ”check_target”]))engine.add_node(Node(”write_alert”, write_alert, deps=[”analyze”, ”check_target”]))engine.add_node(Node(”send_email”, send_email, deps=[”write_report”, ”write_alert”]))# ===== 4. 执行 =====result = engine.run({”target”: 100000, ”week”: ”2026-W30”})print(f”\n🎉 最终输出: {result['send_email']}”)# 预期输出:# --- Layer 0: ['fetch_sales'] ---#   ✅ fetch_sales → {'total': 85000, 'orders': 342}# --- Layer 1: ['analyze', 'check_target'] ---#   ✅ analyze → 销售额85000,订单342单,环比下降12%#   ✅ check_target → fail# --- Layer 2: ['write_report', 'write_alert'] ---#   ✅ write_report → 📊 周报:...#   ✅ write_alert → ⚠️ 预警:...# --- Layer 3: ['send_email'] ---#   ✅ send_email → 邮件已发送: ⚠️ 预警:销售额85000...# 🎉 最终输出: 邮件已发送: ⚠️ 预警:销售额85000...

🎉 这 80 行就是万悟工作流引擎的核心骨架:拓扑排序 → 逐层并行 → 变量传递。万悟在此基础上加了:条件分支(跳过不满足条件的路径)、循环节点、子流程嵌套、重试/熔断、执行持久化、可视化编辑器、触发器系统……但灵魂就是你看到的这个 for layer in layers 循环。


六、总结 & 延伸阅读

本文要点回顾

  1. ✅ 工作流 = 预定义 DAG + 确定性执行,适合步骤固定的企业流程
  2. ✅ 核心算法:Kahn 拓扑排序分层 → 同层并行、跨层串行
  3. ✅ 变量通过 {{node_id.output.field}} 模板在节点间传递
  4. ✅ 条件节点用白名单表达式解析器(非 eval),保障安全
  5. ✅ 重试 + 指数退避处理瞬时故障;fail_fast / continue 两种错误策略
  6. ✅ 工作流与 Agent 互补:固定流程用工作流,开放推理用 Agent

课后作业

  • 在万悟中搭建"智能周报生成"工作流,完成测试运行
  • 配置 Cron 触发器,实现每周一自动执行
  • 添加一个并行节点(同时查两个数据源),观察耗时优化
  • 跑通 Mini Workflow,理解拓扑排序 + 并行执行
  • 思考:如果工作流中某个 LLM 节点输出不稳定(如分类结果偶尔错误),怎么加"校验节点"兜底?

下一篇预告

第 8 篇:多租户 + 权限:企业级 AI 平台的安全基石我们将深入万悟的 RBAC 权限模型、租户隔离机制、API Key 管理,理解为什么"安全"是企业级平台和 Demo 的最大区别。

参考资源

  • 元景万悟 GitHub:https://github.com/UnicomAI/wanwu
  • 工作流引擎独立仓库:https://github.com/UnicomAI/wanwu-workflow
  • Apache Airflow DAG 文档:https://airflow.apache.org/docs/apache-airflow/stable/core-concepts/dags.html
  • Temporal 工作流引擎:https://docs.temporal.io/
  • Kahn 算法:https://en.wikipedia.org/wiki/Topological_sorting#Kahn's_algorithm

📱 关注公众号,追更不迷路

本系列文章首发于微信公众号「农夫三拳有点癫」,每周更新源码拆解与架构实战。

在微信扫描下方二维码即可关注: