工作流引擎:让 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_1.output}} | ||
2.3 万悟支持的节点类型
🔒 代码节点安全——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:什么时候用哪个?
📌 最佳实践:企业级应用中,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.WaitGrouperrCh := make(chan error, len(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) ([][]string, error) {// 构建邻接表和入度表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 [][]stringqueue := []string{}// 入度为 0 的节点是第一层for id, deg := range inDegree {if deg == 0 {queue = append(queue, id)}}for len(queue) > 0 {// 当前层的所有节点layers = append(layers, queue)var nextQueue []stringfor _, 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}} < 100000expression := 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 string) string {// 正则匹配 {{...}} 占位符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 // 默认 3backoff := node.RetryConfig.InitialBackoff // 默认 1svar lastErr errorfor 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 在万悟中搭建
「工作流」→「新建」→ 拖入节点: - 开始节点:定义输入参数 week_start、week_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 分支节点 - 工具节点(邮件) - 工具节点(知识库归档) - 结束节点连线(拖拽) 配置每个节点的参数
▲ 万悟工作流可视化编辑器:拖拽节点、连线、配置参数
4.3 调试执行
点击「测试运行」 输入测试参数: week_start=2026-07-20,week_end=2026-07-26,target=100000观察执行过程:
{week_start: "2026-07-20", ...} | |||
{total: 85000, orders: 342} | |||
false | |||
sent_to: boss@company.com | |||
doc_id: KB-2026-0727 | |||
📌 调试技巧:万悟支持单步执行——点击某个节点可以单独运行,查看它的输入/输出。这对定位"哪个节点出了问题"非常有用。
4.4 配置触发器
0 9 * * 1 | ||
4.5 对比实验
五、自己造一个 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 = idself.func = func# 可调用对象self.deps = deps or []# 依赖的节点 ID 列表# ===== 2. DAG 引擎 =====class WorkflowEngine:def __init__(self):self.nodes = {}self.outputs = {}# node_id → outputdef add_node(self, node):self.nodes[node.id] = nodedef _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] += 1layers, 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] -= 1if in_deg[nb] == 0:next_q.append(nb)queue = next_qif sum(in_deg.values()) > 0:raise ValueError(”Cycle detected!”)return layersdef 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=4) as 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_inputfutures[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循环。
六、总结 & 延伸阅读
本文要点回顾
✅ 工作流 = 预定义 DAG + 确定性执行,适合步骤固定的企业流程 ✅ 核心算法:Kahn 拓扑排序分层 → 同层并行、跨层串行 ✅ 变量通过 {{node_id.output.field}}模板在节点间传递✅ 条件节点用白名单表达式解析器(非 eval),保障安全 ✅ 重试 + 指数退避处理瞬时故障;fail_fast / continue 两种错误策略 ✅ 工作流与 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
📱 关注公众号,追更不迷路
本系列文章首发于微信公众号「农夫三拳有点癫」,每周更新源码拆解与架构实战。
在微信扫描下方二维码即可关注:
夜雨聆风