多 Agent 串行太慢:用 DAG、并发闸门和 Token 预算拆链路
多 Agent 串行太慢:用 DAG、并发闸门和 Token 预算拆链路
多 Agent 链路慢,先别急着换更快的模型。把路由、规划、检索、执行和审查五段的等待时间、输入 Token 与重试次数分别记下来,通常更容易看出问题。
如果所有节点串行运行,还把完整历史透传给下游,延迟会叠加,账单也会包含大量重复上下文。能否改成 DAG,取决于节点之间是否真的有数据依赖;没有依赖的才并行,需要审计的只拿最小输入。
这篇用一组可替换的测试字段说明怎么比较,结论要用自己的模型、请求集和计费规则复跑。
1. 多 Agent 串行链条下的延迟与 Token 瓶颈分析
初始的 Agent 系统架构通常采用串行逻辑。用户提交复杂需求后,系统交给“需求拆解 Agent”,拆出若干子任务;随后依次调用“数据检索 Agent”、“代码生成 Agent”、“安全审计 Agent”以及“文本汇总 Agent”。
用代码逻辑表示,这是一个典型的串行调用:
# 串行调用的示意代码 for step in execution_plan: result = call_llm_agent(step, context_history) context_history.append(result)这段代码是否已经构成瓶颈,不能只凭调用顺序判断。先给每个节点记录开始时间、结束时间、输入输出 Token、重试次数和依赖节点,再用同一批脱敏请求复跑。下面的表格是待采集字段,不预填一组貌似真实的结果:
| Agent 节点 | 耗时分位数 | 输入 Token | 输出 Token | 需要核对的问题 |
|---|---|---|---|---|
| 需求拆解 Agent | 从 Trace 汇总 | 从模型 usage 读取 | 从模型 usage 读取 | 是否包含无关历史 |
| 数据检索 Agent | 从 Trace 汇总 | 从模型 usage 读取 | 从模型 usage 读取 | 是否确实依赖拆解输出 |
| 代码生成 Agent | 从 Trace 汇总 | 从模型 usage 读取 | 从模型 usage 读取 | API 文档是否能按需检索 |
| 安全审计 Agent | 从 Trace 汇总 | 从模型 usage 读取 | 从模型 usage 读取 | 是否收到与审计无关的检索正文 |
| 文本汇总 Agent | 从 Trace 汇总 | 从模型 usage 读取 | 从模型 usage 读取 | 是否重复携带上游全文 |
如果 Trace 显示下游主要在等待前置节点,而且输入 Token 随步骤持续累积,再分别验证两个假设:一是移除虚假依赖,二是按字段裁剪上下文。数据检索与代码生成能否并行,要看二者的输入契约;安全审计需要什么,也应由审计规则决定,不能预先假定它只需要代码。
从底层机制分析,大模型的 Prompt Processing(首字延迟 TTFT)与上下文长度成正比,而 Generation(生成延迟)与输出 Token 数成正比。增加冗余上下文不仅抬高了资金成本,也增加了模型首字响应的卡顿时间。
2. 基于 DAG 有向无环图的执行解耦与上下文隔离
降低延迟的关键在于解耦串行链条,将 Agent 间的依赖关系重构为有向无环图(DAG)。
在重构后的架构设计中,引入“依赖图解析器”与“上下文隔离屏障”。需求拆解完成后,系统分析各个子任务的输入输出契约。只要子任务之间没有数据依赖,就通过异步协程并发调度。同时,每个 Agent 在被触发时,仅接收其特定任务所必需的精简上下文,阻断全局历史的盲目透传。
重构后应缩短请求链路、明确模块边界,并将可复用能力抽到稳定接口中。
若数据检索与代码生成没有依赖,可以在 DAG 中并行;若生成必须引用检索结果,就仍应保留依赖。上下文裁剪后还要回归审计召回率和输出质量,不能只看 Token 下降。
3. 异步并发与 Prompt 裁剪防线代码实现
为了在工程中落地 DAG 并行与 Token 隔离机制,基于 Python 的asyncio和pydantic实现一套轻量级 Agent 调度框架。
代码展示了如何通过异步任务组并发调度无依赖 Agent,并通过上下文提取器精简 Prompt:
import asyncio import time from typing import Dict, Any, List from pydantic import BaseModel, Field class AgentTask(BaseModel): task_id: str agent_type: str dependencies: List[str] = Field(default_factory=list) input_payload: Dict[str, Any] = Field(default_factory=dict) class TokenMetrics(BaseModel): prompt_tokens: int = 0 completion_tokens: int = 0 total_latency_ms: float = 0.0 class AsyncAgentDispatcher: """异步多 Agent 调度器,支持依赖解析与上下文隔离裁剪""" def __init__(self, llm_client): self.llm_client = llm_client self.metrics: Dict[str, TokenMetrics] = {} def _prune_context_for_agent(self, agent_type: str, raw_context: Dict[str, Any]) -> str: """根据 Agent 职责强制裁剪 Prompt,阻断冗余上下文透传""" if agent_type == "security_audit": # 安全审计只需要代码和输入参数,裁剪掉无关的检索文档 code = raw_context.get("generated_code", "") return f"请审计以下代码的安全隐患,仅输出风险项:\n```python\n{code}\n```" elif agent_type == "data_retrieval": query = raw_context.get("user_query", "") return f"根据查询提取关键词并返回检索结果:{query}" elif agent_type == "code_generator": spec = raw_context.get("spec", "") return f"根据以下规格编写 Python 函数:\n{spec}" else: return str(raw_context) async def _execute_single_agent(self, task: AgentTask, context: Dict[str, Any]) -> Dict[str, Any]: start_time = time.perf_counter() pruned_prompt = self._prune_context_for_agent(task.agent_type, context) # 模拟 LLM 异步调用与 Token 统计 prompt_len = len(pruned_prompt) response_text, usage = await self.llm_client.async_generate( agent_type=task.agent_type, prompt=pruned_prompt ) elapsed_ms = (time.perf_counter() - start_time) * 1000 self.metrics[task.task_id] = TokenMetrics( prompt_tokens=usage.get("prompt_tokens", prompt_len // 4), completion_tokens=usage.get("completion_tokens", len(response_text) // 4), total_latency_ms=elapsed_ms ) return {task.task_id: response_text} async def run_dag(self, tasks: List[AgentTask], initial_context: Dict[str, Any]) -> Dict[str, Any]: completed_results: Dict[str, Any] = dict(initial_context) pending_tasks = {t.task_id: t for t in tasks} while pending_tasks: # 筛选出当前依赖已全部就绪的任务 ready_tasks = [ task for task in pending_tasks.values() if all(dep in completed_results for dep in task.dependencies) ] if not ready_tasks: raise RuntimeError("DAG 存在循环依赖或未满足的依赖节点!") # 并行并发执行所有已就绪的 Agent 任务 coroutines = [ self._execute_single_agent(task, completed_results) for task in ready_tasks ] results_list = await asyncio.gather(*coroutines) # 更新上下文并移除已完成任务 for res in results_list: completed_results.update(res) for task in ready_tasks: del pending_tasks[task.task_id] return completed_results工程落地的核心在于_prune_context_for_agent方法。这里用字段白名单代替二次模型摘要,便于审计输入来源。它是否减少延迟和费用,需要用相同请求集比较;若裁剪后任务成功率下降,应补回必要字段,而不是继续压缩。
4. 用同一请求集比较延迟与 Token
准备一组脱敏或合成请求,固定模型版本、并发、缓存状态和重试策略,分别运行串行方案与 DAG 方案。请求数量由环境容量决定,不预设为固定次数。
按任务类型拆分结果,避免一个均值掩盖长上下文或工具超时。表格先保留字段,运行测试后再填值:
| 指标 | 采集来源 | 比较时必须固定的条件 |
|---|---|---|
| P50、P99 与 TTFT | 入口 Trace、模型调用 Trace | 请求集、并发、模型版本、缓存和重试 |
| Prompt / Completion Token | 模型 usage 记录 | Prompt 模板、工具返回与输出质量要求 |
| 单次费用 | 实际 Token 与当前价格表 | 计费区域、缓存折扣和失败重试 |
| 完成吞吐与任务成功率 | 负载工具、任务判定器 | 实例资源、到达率和质量门槛 |
只有输出质量与错误率没有退化,才能把延迟或费用变化归入收益;并行化和裁剪只是待验证的原因。
5. 架构治理与工程防线总结
AI 应用落地不能只等模型自身提速,还要把依赖关系、上下文和预算做成可观测的工程约束。
在搭建多 Agent 协作系统时,建议遵循以下工程规范:
- 阻断无选择的 Context 透传:每个 Agent 仅能读取其执行当前任务所必需的最小数据集合。
- 构建基于 DAG 的异步拓扑:明确 Agent 之间的依赖关系,对于无依赖关系的节点采用异步并发调用。
- 建立 Token 预算与延迟监控告警:线上持续监控各个 Agent 节点的 P99 耗时与 Token 输入输出比例,当 Prompt Token 出现异常增长时,及时触发熔断与排查机制。
将延迟、Token、错误率和任务质量放进同一张监控视图,才能判断这次改图或裁剪究竟有没有改善。
