
一、一个现实问题假设你要 Agent 完成这样一个任务帮我分析这个 Go 项目的性能瓶颈然后优化它提交 PR并在 CI 通过后部署到 staging 环境。这是一个典型的复杂任务包含多个子任务且子任务之间有依赖关系分析代码 → 定位瓶颈 → 修改代码 → 跑测试 → 提 PR → 等 CI → 部署 ↑ CI 必须通过如果 Agent 没有编排能力它可能会还没分析就开始改代码测试没跑完就提 PRCI 还没通过就部署编排Orchestration 就是让 Agent 能合理规划、有序执行复杂任务的能力。二、什么是编排编排是指 Agent 将一个复杂任务分解为多个子任务并按照合理的顺序和依赖关系执行的过程。2.1 编排 vs 简单顺序执行维度简单顺序执行编排任务结构线性一步接一步可能有分支、并行、循环依赖管理隐式靠 prompt 暗示显式用 DAG 定义错误处理从头再来局部重试、跳过、降级状态追踪无有完整的状态机可观测性看日志有可视化的执行图谱2.2 什么时候需要编排任务复杂度是否需要编排示例单步工具调用不需要帮我查一下北京的天气2-3 步线性任务简单编排即可查天气 → 决定带不带伞 → 添加到日历多步且有依赖需要正式编排分析代码 → 改代码 → 测试 → 部署跨系统、多人协作必须编排分析数据 → 生成报告 → 审批 → 发布三、编排的核心DAGDAGDirected Acyclic Graph有向无环图是编排任务的标准数据结构。3.1 DAG 的基本概念┌───────┐ │ 分析代码│ └───┬───┘ │ ┌─────┼─────┐ ▼ ▼ ▼ ┌────┐ ┌────┐ ┌────┐ │定位A│ │定位B│ │定位C│ ← 可以并行 └─┬──┘ └─┬──┘ └──┬─┘ │ │ │ └──────┼───────┘ ▼ ┌───────┐ │ 修改代码│ └───┬───┘ │ ┌───▼───┐ │ 跑测试 │ └───┬───┘ │ ┌───▼───┐ │ 提 PR │ └───┬───┘ │ ┌───▼───┐ │ 等待 CI │ └───┬───┘ │ ┌───▼───┐ │ 部署 │ └───────┘3.2 DAG 的数据结构type Node struct { ID string Task string Status NodeStatus // pending, running, success, failed, skipped DependsOn []string // 依赖的节点 ID } type DAG struct { Nodes map[string]*Node } // 拓扑排序确定节点的执行顺序 func (dag *DAG) TopologicalSort() ([]string, error) { visited : make(map[string]bool) result : make([]string, 0) var dfs func(nodeID string) error dfs func(nodeID string) error { if visited[nodeID] { return nil } visited[nodeID] true node : dag.Nodes[nodeID] for _, depID : range node.DependsOn { if err : dfs(depID); err ! nil { return err } } result append(result, nodeID) return nil } for id : range dag.Nodes { if err : dfs(id); err ! nil { return nil, err } } return result, nil }3.3 DAG 的执行引擎type Orchestrator struct { dag *DAG executor *ToolExecutor } func (o *Orchestrator) Execute(ctx context.Context) error { order, err : o.dag.TopologicalSort() if err ! nil { return fmt.Errorf(拓扑排序失败: %w, err) } results : make(map[string]string) for _, nodeID : range order { node : o.dag.Nodes[nodeID] // 检查依赖是否全部成功 allDepsSuccess : true for _, depID : range node.DependsOn { if o.dag.Nodes[depID].Status ! Success { allDepsSuccess false break } } if !allDepsSuccess { node.Status Skipped continue } // 执行当前节点 node.Status Running result, err : o.executor.Execute(node.Task) if err ! nil { node.Status Failed // 错误处理策略重试、跳过、终止 return fmt.Errorf(节点 %s 执行失败: %w, nodeID, err) } node.Status Success results[nodeID] result } return nil }四、三种编排模式4.1 顺序编排Sequential最简单的模式任务一个接一个执行。A → B → C → D适用场景有严格依赖关系的任务链。示例编译代码 → 跑单元测试 → 打包镜像 → 推送仓库优点简单、可控缺点慢不能利用并行4.2 并行编排Parallel没有依赖关系的任务可以同时执行。┌── B ──┐ A ──┼── C ──┼── E └── D ──┘适用场景多个独立的子任务。示例分析代码结构 ├── 检查内存泄漏 ├── 检查 CPU 热点 ← 三者并行 └── 检查 I/O 瓶颈 汇总报告优点速度快缺点需要处理并发竞争4.3 条件编排Conditional根据前置任务的结果决定下一步走向。A ── 成功 → B ── 成功 → C └── 失败 → D ── 成功 → E └── 失败 → F适用场景需要根据中间结果做决策。示例代码审查 ├── 通过 → 合并到 master → 部署 └── 不通过 → 打回修改 → 重新审查优点灵活、适应性强缺点逻辑复杂测试困难五、LLM 驱动的动态编排上面讲的 DAG 是静态编排——任务结构在运行前就确定了。更高阶的是动态编排Agent 在执行过程中根据当前结果让 LLM 决定下一步做什么。5.1 动态编排的工作流程Step 1: Agent 接收任务 Step 2: LLM 分析任务生成初始计划 Step 3: 执行第一步 Step 4: 观察结果 Step 5: LLM 根据结果调整后续计划 Step 6: 执行下一步 Step 7: 重复 4-6直到任务完成5.2 动态编排 vs 静态编排维度静态编排动态编排计划时机执行前确定边执行边调整确定性高结果可预测低每次可能不同适应性差变更需改 DAG强能应对意外性能快无 LLM 开销慢每一步都要 LLM 决策适用场景流程固定的任务探索性、不确定性任务5.3 混合编排最佳实践实际生产环境中最常用的策略是静态 动态混合1. 用静态 DAG 定义任务的主干流程不变的部分 2. 在每个节点内部用 LLM 动态决策具体执行方式可变的部分示例主干 DAG静态 分析代码 → 修改代码 → 测试 → 部署 分析代码节点内部动态 LLM 决定先看 CPU profile → 再看内存 profile → 最后看锁竞争 如果某个环节发现明显瓶颈提前进入修改阶段六、实战实现一个任务编排引擎package main import ( context fmt sync time ) // 节点状态 type NodeStatus int const ( Pending NodeStatus iota Running Success Failed Skipped ) // 节点定义 type Node struct { ID string Name string DependsOn []string Status NodeStatus Result string Error error Execute func(ctx context.Context) (string, error) } // DAG 定义 type DAG struct { Nodes map[string]*Node } func NewDAG() *DAG { return DAG{Nodes: make(map[string]*Node)} } func (dag *DAG) AddNode(node *Node) { dag.Nodes[node.ID] node } // 获取可并行执行的节点 func (dag *DAG) GetReadyNodes(completed map[string]bool) []*Node { ready : make([]*Node, 0) for _, node : range dag.Nodes { if node.Status ! Pending { continue } allDepsDone : true for _, depID : range node.DependsOn { if !completed[depID] { allDepsDone false break } } if allDepsDone { ready append(ready, node) } } return ready } // 编排引擎 type Engine struct { dag *DAG } func NewEngine(dag *DAG) *Engine { return Engine{dag: dag} } func (e *Engine) Run(ctx context.Context) error { completed : make(map[string]bool) var mu sync.Mutex for { // 检查是否全部完成 allDone : true for _, node : range e.dag.Nodes { if node.Status ! Success node.Status ! Skipped { allDone false break } } if allDone { return nil } // 获取就绪节点 mu.Lock() readyNodes : e.dag.GetReadyNodes(completed) mu.Unlock() if len(readyNodes) 0 { // 有节点未完成但没有就绪节点 → 死锁 return fmt.Errorf(DAG 死锁存在未完成的节点但没有就绪节点) } // 并行执行就绪节点 var wg sync.WaitGroup for _, node : range readyNodes { wg.Add(1) go func(n *Node) { defer wg.Done() n.Status Running fmt.Printf([执行] %s (%s)\n, n.Name, n.ID) result, err : n.Execute(ctx) mu.Lock() if err ! nil { n.Status Failed n.Error err fmt.Printf([失败] %s: %v\n, n.Name, err) } else { n.Status Success n.Result result completed[n.ID] true fmt.Printf([完成] %s: %s\n, n.Name, result) } mu.Unlock() }(node) } wg.Wait() // 如果有节点失败根据策略决定是否继续 // 这里简化遇到失败立即终止 for _, node : range readyNodes { if node.Status Failed { return fmt.Errorf(节点 %s 执行失败: %w, node.Name, node.Error) } } } } func main() { dag : NewDAG() // 定义任务 analyze : Node{ ID: analyze, Name: 分析代码, Execute: func(ctx context.Context) (string, error) { time.Sleep(1 * time.Second) return 发现 3 个性能瓶颈, nil }, } fixCPU : Node{ ID: fix_cpu, Name: 修复 CPU 瓶颈, DependsOn: []string{analyze}, Execute: func(ctx context.Context) (string, error) { time.Sleep(2 * time.Second) return CPU 优化完成, nil }, } fixMemory : Node{ ID: fix_memory, Name: 修复内存泄漏, DependsOn: []string{analyze}, Execute: func(ctx context.Context) (string, error) { time.Sleep(2 * time.Second) return 内存泄漏已修复, nil }, } test : Node{ ID: test, Name: 跑测试, DependsOn: []string{fix_cpu, fix_memory}, Execute: func(ctx context.Context) (string, error) { time.Sleep(1 * time.Second) return 全部测试通过, nil }, } deploy : Node{ ID: deploy, Name: 部署, DependsOn: []string{test}, Execute: func(ctx context.Context) (string, error) { time.Sleep(1 * time.Second) return 部署成功, nil }, } dag.AddNode(analyze) dag.AddNode(fixCPU) dag.AddNode(fixMemory) dag.AddNode(test) dag.AddNode(deploy) engine : NewEngine(dag) err : engine.Run(context.Background()) if err ! nil { fmt.Printf(编排失败: %v\n, err) } else { fmt.Println(全部任务完成) } }输出[执行] 分析代码 (analyze) [完成] 分析代码: 发现 3 个性能瓶颈 [执行] 修复 CPU 瓶颈 (fix_cpu) [执行] 修复内存泄漏 (fix_memory) [完成] 修复内存泄漏: 内存泄漏已修复 [完成] 修复 CPU 瓶颈: CPU 优化完成 [执行] 跑测试 (test) [完成] 跑测试: 全部测试通过 [执行] 部署 (deploy) [完成] 部署: 部署成功 全部任务完成注意fix_cpu和fix_memory是并行执行的因为它们都只依赖analyze彼此没有依赖关系。七、编排的最佳实践7.1 超时控制每个节点都应该有超时限制防止某个任务卡死整个流程ctx, cancel : context.WithTimeout(parentCtx, 30*time.Second) defer cancel() result, err : node.Execute(ctx)7.2 重试策略对于可重试的失败网络超时、限流自动重试type RetryStrategy struct { MaxRetries int Backoff time.Duration } func executeWithRetry(ctx context.Context, node *Node, strategy RetryStrategy) (string, error) { var lastErr error for i : 0; i strategy.MaxRetries; i { result, err : node.Execute(ctx) if err nil { return result, nil } lastErr err time.Sleep(strategy.Backoff * time.Duration(1 i)) // 指数退避 } return , fmt.Errorf(重试 %d 次后仍然失败: %w, strategy.MaxRetries, lastErr) }7.3 可见性每个节点的执行状态应该对外暴露方便监控和调试type NodeExecutionEvent struct { NodeID string Status NodeStatus Timestamp time.Time Duration time.Duration Result string Error string }7.4 人工介入点在关键节点设置需要人工确认的闸门type GateNode struct { Node RequiresApproval bool ApprovedBy string } // 执行到闸门节点时暂停等待人工确认后再继续八、课后实践任务用 Go 实现一个支持并行执行和错误处理的 DAG 编排引擎。要求支持节点的依赖关系定义支持并行执行无依赖的节点支持节点的超时控制每个节点最多执行 5 秒支持节点的重试机制失败后最多重试 2 次输出完整的执行日志谁在什么时候执行了什么结果如何进阶挑战实现条件分支根据前置节点的结果决定走哪条路径实现人工审批节点执行到该节点时暂停等待外部确认信号把编排引擎和 zz365.top 的工具集成——比如用 Crontab 定时触发 DAG 执行延伸思考如果 DAG 中有 100 个节点你怎么可视化展示执行进度如果某个节点执行了 10 分钟还没结束怎么判断它是正常执行还是卡死了多个 Agent 共享同一个 DAG 时怎么避免资源竞争九、延伸阅读论文《Plan-and-Solve Prompting: Improving Zero-Shot Chain-of-Thought Reasoning by Large Language Models》2023- 任务规划的早期研究论文《TaskMatrix: Connecting Foundation Models with Diverse APIs》2023- 任务编排与工具调用的结合开源项目Temporal / Airflow - 生产级的任务编排引擎参考其 DAG 设计思路开源项目Dagger - Go 实现的 CI/CD 编排引擎十、下一讲预告第6讲沙箱与安全——让 Agent 安全地执行代码我们会深入探讨Agent 执行代码时怎么保证安全Docker 沙箱、gVisor、WASM 沙箱的选型和实现以及怎么防止 Prompt 注入攻击。开发之余的小工具推荐处理 Base64、JWT 解析、JSON 格式化、Crontab 计算、PDF 合并压缩这些碎片需求我常用一个纯前端本地工具箱zz365.top。所有计算在浏览器完成文件不上服务器关页即清。免费、无登录、无广告适合开发者当常驻标签页。