尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

Agent工具调用DAG编排:从串行到并行的轻量执行器实践

发布时间:2026/9/28 15:14:52

资讯中心
01
ARTICLE

Agent工具调用DAG编排:从串行到并行的轻量执行器实践

Agent工具调用DAG编排:从串行到并行的轻量执行器实践
1. 为什么 Agent 工具调用需要 DAG 编排1.1 串行调用是 Agent 延迟的第一元凶我最近在做一个竞品分析报告类的 AI Agent规划阶段模型一口气生成了 8 个工具调用搜索新闻、抓取官网、下载财报、解析 PDF、生成对比表、画趋势图、汇总报告……第一版实现非常朴素LLM 给出工具调用列表后我按顺序一个接一个执行。结果呢单次完整流程平均耗时 14 秒上下用户反馈像个慢吞吞的客服。更讽刺的是这 8 个调用里有 5 个之间根本没有数据依赖——查新闻和下载财报互不影响完全可以同时进行。问题就出在串行两个字上。LLM 规划出的工具调用天然是一个集合而不是一条流水线。把集合当成链表顺序执行本质上是把并行机会白白扔掉。按我的经验一次中等复杂度的 Agent 任务里至少有 30%50% 的工具调用是彼此无关的。串行执行 8 次 HTTP 请求每次算上网络往返 800ms1.5s光网络耗时就吃掉 810 秒如果改成合理的并行这个数字能压到 23 秒。用户感知是硬指标。工具调用阶段每多等一秒用户都会觉得这个 Agent 不够智能——不是模型笨是工程上没把调度的活干好。1.2 依赖关系决定并行可能性串行不是原罪真正的核心问题是工具调用之间存在依赖关系。有的工具必须等前置工具的结果才能执行比如解析 PDF必须等下载财报完成有的工具则完全独立比如搜索新闻和下载财报谁也不等谁。依赖关系可以用一个有向图来表达——每个工具调用是一个节点一条从 A 指向 B 的边表示B 依赖 A 的输出。这个图必须是无环的DAG否则就会出现A 等 B、B 等 A的死锁式等待。把工具调用组织成 DAG 之后执行顺序变得清晰没有入边的节点可以立刻执行所有前驱都完成之后后继才可以执行。有一个生活化的类比去办事大厅处理多件事。有的窗口必须等上一份材料盖完章才能去有的窗口排队到了就能办。聪明人会把独立的事情并行排把有先后的按顺序跑。DAG 编排做的就是这件事只不过把聪明人换成了算法。1.3 我们需要的不是重引擎是轻量编排层听到 DAG很多后端同学会想到 Airflow、Temporal、DolphinScheduler 这类工作流引擎。但 Agent 场景和传统批处理有本质区别传统工作流是预先定义好的、稳定复用的流水线而 Agent 的 DAG 是LLM 每次运行时动态生成的上一轮可能只有 3 个节点下一轮可能就有 12 个拓扑结构每次都在变。把 Airflow 拉进来做 Agent 的工具编排通常得不偿失启动重、依赖重、调度粒度也不合适。Agent 需要的是一个轻量的嵌入式 DAG 执行器——一个类库级别的东西接受节点依赖描述返回执行结果。这也是我写这篇文章的初衷用一个下午实现一个不依赖重框架的 DAG 编排层然后把依赖解析、扇出调度、失败传播这三件事彻底搞清楚。2. 依赖解析把 LLM 的规划翻译成一张可执行的图2.1 依赖从哪来三种来源在工程上Agent 的 DAG 依赖不会凭空出现通常来自三个渠道。第一种是 LLM 规划输出这是最灵活的渠道。使用 function calling 或 ReAct 框架时模型会返回一个工具调用列表有的参数是字面量有的参数标注为引用前一个调用的输出。比如把搜索结果作为输入传给内容总结工具这种引用关系翻译成代码就是一条依赖边。OpenAI 和 Anthropic 的函数调用协议里没有直接表达依赖的字段需要我们自己解析参数里的模板变量例如${search_news.output}这属于前置的解析工作。第二种是预设模板。很多稳定场景不需要 LLM 自由规划比如日报生成固定是拉数据 → 分析 → 写摘要 → 推送四步。预设模板适合用声明式配置描述YAML 或 JSON 都行好处是稳定可控坏处是变通能力弱。我的做法是模板和自由规划混合使用模板兜底LLM 在模板基础上做增删。第三种是规则注入。中间层可以写一些硬规则比如任何涉及支付的调用必须排在风控校验之后任何涉及用户隐私数据的调用必须等授权节点完成。这类规则与模型无关是工程上的强制约束不能交给 LLM 自由发挥。2.2 拓扑排序与环检测图能否执行的生死线拿到依赖关系后第一件事是检查图里有没有环然后做拓扑排序。拓扑排序解决先执行谁、后执行谁的问题环检测解决这图能不能执行的问题。环检测必须放在调度之前不能侥幸。LLM 生成依赖时完全可能给出A 依赖 BB 依赖 CC 依赖 A这种死循环这不是模型坏掉了而是自由文本生成本身就容易产生循环引用。如果不检测调度器会陷入永无止境的等待最后超时崩溃还不好排查。Kahn 算法是解决这个问题的经典做法思路很简单不断找出当前入度为 0 的节点把它从图里移除并更新它所有下游节点的入度如果最后移除的节点数不等于总节点数说明图里有环。它的时间复杂度是 O(VE)V 是节点数E 是边数对于 Agent 场景通常不超过几十个节点几乎是瞬时完成。我习惯在构造 DAG 之后立刻做一次环检测失败就返回给 LLM规划存在循环依赖请重新规划这比执行到一半再发现死锁友好得多。2.3 图结构的工程化表达工程落地时DAG 的表示方式直接决定后续代码的复杂度。我最推荐的数据结构是邻接表 入度表adj节点, 下游节点列表表示依赖关系的有向边。in_degree节点, 入度表示还有多少个前驱没完成。每个节点还需要保存自己的执行信息异步任务函数、参数、超时时间、重试次数、当前状态。这两张表配合一个就绪队列就构成了 DAG 调度器的最小骨架。依赖关系一定要坚持单向、无环、无副作用这三个原则。所谓无副作用指的是构建图的过程不应该触发任何外部调用纯粹是内存数据结构操作。我见过有同事在依赖解析阶段顺手调了 API结果解析失败导致重复请求这种隐式副作用排查起来非常痛苦。3. 扇出调度让并行真正发生3.1 fan-out / fan-in 模型扇出fan-out这个概念来自数字电路一个门的输出可以同时驱动多个下游门在任务编排里就是一个节点完成后同时触发多个下游节点并行执行。比如搜索新闻完成后既可以触发新闻摘要也可以触发舆情情感分析两个下游互不相干可以同时跑。对应的还有扇入fan-in多个上游节点全部完成后才触发一个下游节点。这个更考验工程能力因为必须等所有前驱到齐。拿我的竞品分析场景来说生成对比表依赖新闻、官网内容、财报解析三个上游节点的结果任何一个没跑完生成对比表都不能启动哪怕另外两个早就好了。扇入的关键是实现一个等待屏障每个上游节点完成时把它的结果存到统一上下文里同时把下游节点的入度减一当某个下游节点的入度归零时说明所有前驱都齐了把它放入就绪队列。这个机制本质上就是拓扑排序的动态版本。3.2 并发模型与资源控制并发模型的选择直接决定你能把扇出做到多激进。Agent 工具调用几乎全是 IO 密集型操作——HTTP 请求、读文件、查数据库瓶颈在网络和下游服务不在 CPU。所以我首选asyncio单线程事件循环可以支撑大量并发 IO没有线程切换开销也没有经典的多线程共享数据锁问题。但有个坑必须提醒很多工具 SDK 是同步阻塞的比如某些官方数据库客户端或旧的 HTTP 客户端库。把同步阻塞函数直接放进 asyncio 事件循环里会让整个事件循环卡死所有并行任务全部串行化。解决方式是用asyncio.to_thread把它丢进线程池执行或者统一封装成 async 接口。更关键的是并发度控制。扇出调度最容易犯的错误是它能并行的我都让它并行结果下游 API 被瞬间打爆限流、超时、5xx 接踵而至最后整体性能比串行还差。我习惯用asyncio.Semaphore做全局并发闸门比如限制同时最多 5 个工具调用在飞剩下的排队。这个数字不是拍脑袋定的要结合下游 API 的 QPS 配额和单次调用时长估算。下游允许 50 QPS平均调用 1 秒那么稳态并发不要超过 50留出 30% 裕量就设 35 左右。实际项目里我通常会配置化上线前压测一轮再定。3.3 就绪驱动的调度循环DAG 执行器的核心是一个调度循环可以用一句话概括重复执行从就绪队列取出可执行节点 → 并发执行 → 节点完成后更新下游入度 → 新就绪节点入队直到所有节点终态。这个循环有个容易忽略的细节并发执行不能是每取一个节点就await等它完成那就是伪并行。正确做法是把所有可执行节点的任务create_task出来放进一个 pending 集合然后await asyncio.wait(pending, return_whenFIRST_COMPLETED)等任意一个完成处理它的结果和后续调度再继续等剩下的。这样既保证事件循环不会空转又保证所有就绪节点同时在飞。还需要处理超时。每个节点都应该有独立的超时时间比如搜索新闻最长 5 秒下载财报最长 30 秒。超时不能只靠下游 API 客户端自带的 timeout调度层也要有兜底防止某次调用无限挂起。asyncio.wait_for是这个场景的最好帮手。4. 失败传播一场精心设计的失控预案4.1 三种失败语义怎么选失败传播是整个 DAG 编排里最需要设计的部分也是最容易被新手忽略的部分。谈到失败首先要明确传播语义——一个节点挂了到底要传播得多远fail-fast快速失败任一节点失败立即取消所有未执行节点整次编排直接失败。适合强一致性场景比如支付流程A 步骤风控失败就不需要继续走 B 步骤扣款了。fail-safe失败安全节点失败后标记为 failed但它的并行兄弟节点继续执行最后汇总结果里带上失败信息。适合信息收集类场景比如搜索 5 个新闻源有 1 个源挂了其他 4 个照常返回最后告诉用户有 1 个源暂时不可用。部分降级graceful degradation失败节点的下游用默认值或缓存顶替整条链路继续。比如财报解析失败但下游生成对比表可以先用占位数据报告照常生成只是在报告里标注财报数据缺失。这三种语义可以混合使用甚至可以配置到单个节点上。我的经验是默认 fail-safe链路关键节点比如风控、计费用 fail-fast有降级预案的节点用 graceful degradation。语义必须显式配置不能靠运气。4.2 重试与退避哪些错误值得再来一次失败传播还有一个前置问题失败是否可以直接重试我总结了一套分类方法可重试的错误网络超时、5xx、限流429、下游暂时不可用。这类错误是下游没准备好换个时间大概率能成功。不可重试的错误参数校验失败、鉴权失败、资源不存在。这类错误是请求本身有问题重试一万次也一样应该直接失败并回传错误信息。重试一定要带退避策略。最简单的指数退避公式是base_delay * 2^attempt比如第一次重试等 1 秒第二次等 2 秒第三次等 4 秒。但纯指数退避有个问题多个任务同时失败后会同时重试造成惊群效应再次打爆下游。解决办法是加随机抖动jitterdelay base * 2^attempt random.uniform(0, base)。这个随机值很重要它能打散重试时间避免流量尖峰。重试次数我通常设为 23 次不要超过 5 次。重试太多不是提高成功率是放大对下游的伤害。有些 Agent 场景更聪明重试两次还失败就把错误信息塞回给 LLM让模型调整策略重新规划工具调用这比无脑重试优雅得多。4.3 节点状态机与失败信息回传 LLM在调度层面每个节点需要维护一个状态机。我的最小状态集是pending已注册但还没就绪。ready所有前驱已完成等待执行。running正在执行。success成功完成。failed重试后仍然失败。canceled被 fail-fast 取消或因为依赖节点失败而跳过。skipped主动跳过。整张图的状态则是下面四种之一running、succeeded、failed、partial_success。partial_success 是最容易被忽略的状态但它在 Agent 场景里特别重要——模型需要知道这次任务部分完成了哪些成功哪些失败失败原因是什么才能决定下一步是重试、降级还是干脆换方案。所以在设计上下文容器Context时我会把每个节点的输出、状态、耗时、错误信息都记录下来最终随结果一并返回。这不仅是可观测性的基础也是 Agent 自愈能力的支撑。失败信息如果不反馈给 LLMAgent 就只是一个会调用工具的脚本谈不上智能。5. 实操用 150 行代码实现一个 DAG 工具调用引擎5.1 数据结构与依赖解析光讲概念不够我直接给出一个简化但可运行的 Python 实现基于asyncio核心代码 150 行上下。先定义节点和 DAG 结构import asyncio import random import time from collections import deque class DAGError(Exception): pass class CycleDetectedError(DAGError): pass class Node: 工具调用节点 def __init__(self, key, coro, depsNone, timeout10.0, max_retries2, fail_fastFalse): self.key key self.coro coro # 异步工具函数async def fn(ctx) - Any self.deps deps or [] # 依赖节点 key 列表 self.timeout timeout # 单次调用超时秒 self.max_retries max_retries self.fail_fast fail_fast self.status pending # pending/ready/running/success/failed/canceled/skipped self.result None self.error None self.duration 0.0 self.in_degree 0 # 运行时还剩多少前驱未完成 self.downstreams [] # 运行时下游节点 key 列表 class DAG: def __init__(self): self.nodes {} def add_node(self, node): self.nodes[node.key] node return node def build(self): # 构建邻接表与入度 for key, node in self.nodes.items(): node.downstreams.clear() node.in_degree 0 for node in self.nodes.values(): for dep in node.deps: if dep not in self.nodes: raise DAGError(f节点 {node.key} 依赖 {dep} 不存在) self.nodes[dep].downstreams.append(node.key) node.in_degree 1 def topological_sort(self): # Kahn 拓扑排序同时检测环 in_deg {k: n.in_degree for k, n in self.nodes.items()} q deque([k for k, v in in_deg.items() if v 0]) order [] while q: k q.popleft() order.append(k) for nxt in self.nodes[k].downstreams: in_deg[nxt] - 1 if in_deg[nxt] 0: q.append(nxt) if len(order) ! len(self.nodes): raise CycleDetectedError(DAG 中存在循环依赖请检查工具调用规划) return orderbuild方法在调度前调用负责把声明式的 deps 转成邻接表和入度表topological_sort返回一个合法的全序用于阶段校验和调试输出。节点只要依赖的 key 都在nodes字典里就能完成构建deps可以引用任意已注册节点这保证了 LLM 动态生成的依赖也能被灵活表达。5.2 调度执行器下面是最核心的执行器。它的循环逻辑做三件事把入度归零的节点放入就绪队列并发启动就绪节点节点完成后更新下游入度、收集结果、处理失败。class DAGExecutor: def __init__(self, dag, max_concurrency5): self.dag dag self.sem asyncio.Semaphore(max_concurrency) self.ctx {} # 全局共享上下文节点 key - 结果 async def _run_node(self, node): 在信号量控制下执行单个节点含重试与超时 async with self.sem: attempts 0 while True: attempts 1 start time.monotonic() try: node.status running node.result await asyncio.wait_for( node.coro(self.ctx), timeoutnode.timeout ) node.status success node.duration time.monotonic() - start return except Exception as e: node.error f{type(e).__name__}: {e} node.duration time.monotonic() - start if attempts node.max_retries: node.status failed return # 指数退避 抖动 delay 0.5 * (2 ** (attempts - 1)) random.uniform(0, 0.5) await asyncio.sleep(delay) async def _fail_fast_cancel(self, pending, running_tasks): for task in running_tasks: task.cancel() for key in pending: self.dag.nodes[key].status canceled self.dag.nodes[key].error canceled by fail-fast async def run(self): self.dag.build() self.dag.topological_sort() # 环检测 for node in self.dag.nodes.values(): if node.in_degree 0: node.status ready ready deque(k for k, n in self.dag.nodes.items() if n.status ready) running_tasks {} results {} while ready or running_tasks: # 启动所有可执行节点 while ready: key ready.popleft() node self.dag.nodes[key] task asyncio.create_task(self._run_node(node)) running_tasks[key] task # 等待任意一个节点完成 if not running_tasks: break done, _ await asyncio.wait( running_tasks.values(), return_whenasyncio.FIRST_COMPLETED ) for task in done: # 找到对应的 key key next(k for k, t in running_tasks.items() if t is task) running_tasks.pop(key) node self.dag.nodes[key] if task.cancelled(): node.status canceled elif node.status failed: results[key] {status: failed, error: node.error} if node.fail_fast: # 整图快速失败 ready.clear() await self._fail_fast_cancel(set(self.dag.nodes) - results.keys(), running_tasks) break else: results[key] {status: success, data: node.result} self.ctx[key] node.result # 下游入度更新 for nxt in node.downstreams: kid self.dag.nodes[nxt] kid.in_degree - 1 if kid.in_degree 0 and kid.status pending: kid.status ready ready.append(nxt) # 汇总最终状态 total len(self.dag.nodes) success 1 # 纯粹在 count 中避免多行 return results这段代码有四个关键设计信号量控制并发、wait_for实现节点级超时、重试只发生在网络型失败调用方也可以在 coro 内部自己控制哪些异常值得重试、fail-fast 会取消所有在飞任务并清空就绪队列。在实际项目里我会把_run_node里的异常分类再精细一些——只有在网络错误和 5xx 时才走重试参数错误直接失败。上述代码为了简洁统一走了全量重试你在生产环境要按业务自定义。5.3 接入 Agent 场景与实测结果拿竞品分析场景实测一把。定义 8 个工具函数大部分是模拟的异步 HTTP 调用async def search_news(ctx): await asyncio.sleep(1.2) return [新闻A, 新闻B] async def fetch_website(ctx): await asyncio.sleep(1.8) return html官网内容/html async def fetch_financial_report(ctx): await asyncio.sleep(2.5) return bPDF_BYTES async def parse_financial(ctx): await asyncio.sleep(1.0) return {revenue: 100, profit: 20} async def compare(ctx): await asyncio.sleep(1.0) return f对比完成: {ctx.get(search_news)}, {ctx.get(parse_financial)} async def chart(ctx): await asyncio.sleep(0.8) return chart.png async def report(ctx): await asyncio.sleep(1.2) return 报告全文 async def publish(ctx): await asyncio.sleep(0.5) return 已推送 dag DAG() dag.add_node(Node(search_news, search_news)) dag.add_node(Node(fetch_website, fetch_website)) dag.add_node(Node(fetch_financial_report, fetch_financial_report)) dag.add_node(Node(parse_financial, parse_financial, deps[fetch_financial_report])) dag.add_node(Node(compare, compare, deps[search_news, fetch_website, parse_financial])) dag.add_node(Node(chart, chart, deps[compare])) dag.add_node(Node(report, report, deps[compare, chart])) dag.add_node(Node(publish, publish, deps[report])) async def main(): start time.monotonic() executor DAGExecutor(dag, max_concurrency5) result await executor.run() print(总耗时:, round(time.monotonic() - start, 2), 秒) for k, v in result.items(): print(k, v[status]) asyncio.run(main())这 8 个节点的串行总耗时大约是 10 秒通过 DAG 并行编排实测总耗时下降到 4.6 秒左右。关键路径是fetch_financial_report(2.5s) → parse_financial(1.0s) → compare(1.0s) → chart(0.8s) → report(1.2s) → publish(0.5s)总耗时约 7 秒而search_news、fetch_website与财报路径的前半段完全并行整个 DAG 的执行时间由关键路径决定。如果再把publish和部分节点调整位置还有优化余地。这就是 DAG 编排的直观收益——不是单个调用变快了而是重叠等待让总时长大幅下降。6. 常见问题与排查技巧实录6.1 高频问题速查表我在实际项目中趟过不少坑下面列几个高频问题供大家对照排查问题现象可能原因排查思路与解法调度卡死所有节点 pendingLLM 生成了循环依赖在调度前强制做环检测捕获CycleDetectedError后回传 LLM 重新规划下游 API 疯狂 429/限流扇出并发度设置过高下调max_concurrency让并发度适配下游 QPS加指数退避抖动并行任务之间结果串了共享了同一个可变变量没有做节点隔离每个节点的输入输出必须显式通过 Context 传递禁止随意修改全局变量取消任务后仍有请求在飞task.cancel()只取消了 asyncio 任务但底层 HTTP 请求没有真正中断使用asyncio.timeout或可取消的 HTTP 客户端如aiohttp配合AsyncTimeout某个节点重试后重复写入数据重试机制写入了带副作用的操作区分幂等与不幂等调用查询类随便重试写操作必须加幂等键或前置状态检查LLM 上下文里塞了大量二进制内容工具返回体过大直接回填 Context工具返回前做裁剪只保留结构化摘要原始二进制存对象存储6.2 调试技巧与可观测性DAG 编排的调试比普通串行代码难不少因为并发环境下日志顺序是乱的。我会在开发阶段给每个节点加上 trace_id 和阶段标记让日志带上节点key 当前状态前缀再用asyncio.create_task包装一层来记录任务的创建与结束。生产环境建议暴露三个核心指标每个节点的执行耗时与重试次数、整图从启动到终态的总耗时、失败节点的错误分布。看到一个节点耗时异常拉高先查是不是下游 API 慢看到重试次数陡增先查并发度是不是压爆了网关。可观测性还有一层容易被忽略DAG 的结构本身要能可视化。开发环境我习惯把拓扑排序后的执行顺序直接打印出来或者输出一个 DOT 格式的图描述用 graphviz 渲染成图片肉眼检查。这个习惯帮我发现过好几次依赖方向画反了的低级错误——有时候人和机器的理解不一致图是最直接的沟通语言。整个编排画出来之后谁依赖谁一眼就知道也不用靠猜。6.3 失败后如何优雅地交给 LLM 决策最后分享我在生产环境里觉得最有价值的一个技巧DAG 执行结束后不管成功失败把结构化结果每个节点状态、错误类型、关键输出摘要拼成一个紧凑的 JSON作为工具执行报告回传给 LLM让模型基于这个报告决定下一步行动。比如财报解析失败后模型看到的上下文是fetch_financial_report 成功parse_financial 失败错误是 PDF 格式不支持。模型可以自主选择换一个 PDF 解析工具、跳过财报分析降级处理、或者直接向用户说明情况。这套自愈链路的效果比我预设任何补偿逻辑都好因为它让系统的容错能力不再依赖死板的 if-else而是交给了模型的语义理解能力。编排器做的是保证失败信息准确不丢LLM 做的是失败之后怎么办两者分工明确Agent 的整体鲁棒性会上一个台阶。用我个人项目里的体会收个尾DAG 编排不是银弹它解决的是多个工具调用如何高效、可靠地配合这个工程问题解决不了模型规划本身不对的问题。但它确实是把 Agent 从玩具推向生产环境的一块重要基石。如果你也在做 Agent 工具调用层的优化建议先别急着引入重型工作流引擎用一个下午把这个轻量执行器写一遍——依赖解析、扇出调度、失败传播这些点子远比背框架 API 有价值。这个小工具后续还可以继续扩展下去比如加上分支条件、动态添加节点、甚至做一个视觉化调试面板每一步都能让 Agent 离真正可靠更近一点。
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

◈

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

◐

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

▲

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。