三亩地 三亩地SAN MU DI · CODE DIARY
ARTICLE DETAIL

日记详情

真实记录编程学习的某一天,欢迎挑你感兴趣的翻一翻。

多 Agent 串行太慢:用 DAG、并发闸门和 Token 预算拆链路

多 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 的asynciopydantic实现一套轻量级 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 协作系统时,建议遵循以下工程规范:

  1. 阻断无选择的 Context 透传:每个 Agent 仅能读取其执行当前任务所必需的最小数据集合。
  2. 构建基于 DAG 的异步拓扑:明确 Agent 之间的依赖关系,对于无依赖关系的节点采用异步并发调用。
  3. 建立 Token 预算与延迟监控告警:线上持续监控各个 Agent 节点的 P99 耗时与 Token 输入输出比例,当 Prompt Token 出现异常增长时,及时触发熔断与排查机制。

将延迟、Token、错误率和任务质量放进同一张监控视图,才能判断这次改图或裁剪究竟有没有改善。

← 返回列表