AI Agent 工具调度与并发执行:从同步循环到 DAG 异步执行引擎
当 AI Agent 的"大脑"能够同时使用多个工具时,如何让工具调用从串行等待变为智能并行,是决定 Agent 响应延迟和吞吐量的核心工程问题。
一、工具调用的瓶颈:为什么同步循环正在拖垮 Agent
1.1 当前 Agent 运行时的典型困境
现代 AI Agent 在处理复杂任务时,往往需要连续调用多个工具:搜索数据库、读取文件、调用 API、执行代码……在绝大多数开源框架中(LangChain、CrewAI、AutoGen 的早期实现),工具调用默认采用同步循环模式——Agent 发出一个 LLM 请求,根据返回的 tool_call 执行单个工具,将结果追加回对话历史,然后再发出下一个 LLM 请求。
[LLM 请求] → [等待响应] → [解析 tool_call] → [执行工具A] → [等待结果] → [追加上下文] → [LLM 请求] → ...
这种模式下,假设一个任务需要调用 5 个工具,每个工具平均耗时 200ms,LLM 推理单次耗时 1.5s,总延迟为:
$$T_{sync} = \sum_{i=1}^{n}(T_{LLM} + T_{tool_i}) = 5 \times (1.5s + 0.2s) = 8.5秒$$
而理论上最优的并行模式——当多个工具之间不存在数据依赖时同时发起——则可能将延迟压缩到:
$$T_{parallel} = T_{LLM} + \max(T_{tool_1}, ..., T_{tool_n}) + T_{LLM} = 1.5s + 0.2s + 1.5s = 3.2秒$$
理论上 2.6 倍的延迟差距,在实际生产环境中往往更为悬殊——当工具涉及网络 I/O、数据库查询和外部 API 时,同步阻塞的时间会被进一步放大。
1.2 问题的本质
工具调度的并发性问题本质上是一个任务图调度问题。我们需要回答三个核心问题:
- 依赖分析:哪些工具调用的输出是另一个工具的输入?
- 并行分组:哪一批工具可以安全地并行执行?
- 背压控制:当并行工具数量超过资源限制时如何排队?
这些问题在分布式计算和编译器优化领域早有成熟方案,但 AI Agent 工程界直到最近一年才开始系统性地引入这些技术。
二、从 ReAct 到 Plan-and-Execute:范式演进的脉络
2.1 ReAct 模式的天花板
ReAct(Reasoning + Acting)是当前 Agent 架构的基础范式,其核心思想是交替执行"思考"和"动作"。这种模式的优点是逻辑清晰、可追踪,缺点是天然串行——因为每一步推理都依赖上一步的观察结果。
# 典型 ReAct 同步循环伪代码
async def react_agent_loop(user_query: str):
messages = [SystemMessage(content=SYSTEM_PROMPT), HumanMessage(content=user_query)]
while True:
response = await llm.ainvoke(messages) # 等待 LLM 响应
messages.append(response)
if not response.tool_calls:
return response.content # 无工具调用则结束
for call in response.tool_calls: # 串行执行工具
tool_result = await execute_tool(call.name, call.args)
messages.append(ToolMessage(content=str(tool_result), tool_call_id=call.id))
2.2 Plan-and-Execute 的解耦思路
Plan-and-Execute 模式将"计划"和"执行"拆分为两个独立阶段:
- Planner:生成完整的多步骤执行计划(通常是一个步骤列表)
- Executor:根据计划,尽可能并行地执行无依赖的步骤
Planner LLM 调用:
"用户问北京天气和股票行情"
→ Plan: [Step1: 搜索天气, Step2: 获取股票] (两个步骤无依赖)
Executor 并行调度:
├── TaskPool.submit(weather_search, "北京")
└── TaskPool.submit(stock_query, "沪指")
两个任务并行执行,不需要等待另一个完成
这种模式天然支持并行,但面临两个挑战:计划阶段的 LLM 可能生成不合理的步骤划分,以及执行阶段的错误恢复机制较为复杂。
2.3 工业界的折中方案
在实际工程中,我们看到三种主流路线:
| 方案 | 代表框架 | 并行能力 | 灵活性 |
|---|---|---|---|
| 纯 ReAct | LangChain Agent | 无 | 高(动态决策) |
| Plan-then-Execute | LangGraph, CrewAI | 中等 | 中(需离线规划) |
| 混合模式 | Anthropic Claude Agent SDK | 高 | 高(运行时决定并行) |
三、DAG 执行引擎:将编译器优化思想引入 Agent 调度
3.1 将工具调用抽象为有向无环图
Anthropic 在 Claude Agent SDK 中的设计给了我们关键启发:当 LLM 返回多个 tool_call 时,它们的参数中可能包含 __reference__ 占位符,指向之前工具的输出。这种引用关系天然构成一个 DAG。
// 工具调用的 DAG 节点定义
#[derive(Debug, Clone)]
struct ToolNode {
id: ToolCallId,
name: String,
args: Value,
// 依赖的其他工具输出 ID 列表
dependencies: HashSet<ToolCallId>,
// 执行状态
state: ExecutionState,
}
#[derive(Debug, Clone)]
enum ExecutionState {
Pending, // 等待依赖完成
Ready, // 所有依赖已满足,可以执行
Running, // 正在执行
Completed(Value), // 执行完成,附带结果
Failed(String), // 执行失败
}
// DAG 结构
struct ToolDAG {
nodes: HashMap<ToolCallId, ToolNode>,
// 邻接表:当前节点 → 依赖它的节点列表
reverse_edges: HashMap<ToolCallId, HashSet<ToolCallId>>,
}
3.2 拓扑排序与并行调度
构建 DAG 后,我们可以通过 Kahn 算法进行拓扑排序,然后按层级并行调度:
impl ToolDAG {
/// 找出当前所有 Ready 状态的节点
fn find_ready_nodes(&self) -> Vec<ToolCallId> {
self.nodes.iter()
.filter(|(_, node)| matches!(node.state, ExecutionState::Ready))
.map(|(id, _)| id.clone())
.collect()
}
/// 标记节点完成,将其下游依赖的依赖计数减一
fn complete_node(&mut self, id: ToolCallId, result: Value) {
self.nodes.get_mut(&id).unwrap().state = ExecutionState::Completed(result);
if let Some(dependents) = self.reverse_edges.get(&id) {
for dep_id in dependents {
let dep_node = self.nodes.get_mut(dep_id).unwrap();
dep_node.dependencies.remove(&id);
if dep_node.dependencies.is_empty() {
dep_node.state = ExecutionState::Ready;
}
}
}
}
}
3.3 异步执行引擎的完整实现
use tokio::sync::{mpsc, Semaphore};
use tokio::task::JoinSet;
use std::sync::Arc;
struct AsyncExecutor {
tool_registry: Arc<dyn ToolRegistry>,
concurrency_limit: Arc<Semaphore>,
max_retries: usize,
}
impl AsyncExecutor {
async fn execute_dag(&self, mut dag: ToolDAG) -> Result<DAGResult, ExecutorError> {
let (tx, mut rx) = mpsc::channel::<(ToolCallId, Result<Value, ToolError>)>(256);
let mut join_set = JoinSet::new();
loop {
// 找出所有就绪节点
let ready_nodes = dag.find_ready_nodes();
if ready_nodes.is_empty() {
// 没有就绪节点也没有运行中的任务 → 全部完成
if join_set.is_empty() {
break;
}
} else {
// 批量提交就绪节点
for node_id in ready_nodes {
let node = dag.nodes.get(&node_id).unwrap().clone();
let permit = self.concurrency_limit.clone().acquire_owned().await?;
let tx = tx.clone();
let registry = self.tool_registry.clone();
// 替换参数中的引用占位符
let resolved_args = self.resolve_references(&node.args, &dag)?;
join_set.spawn(async move {
let _permit = permit; // 持有信号量许可直到任务完成
let result = registry.execute(&node.name, resolved_args).await;
let _ = tx.send((node_id, result)).await;
});
dag.nodes.get_mut(&node_id).unwrap().state = ExecutionState::Running;
}
}
// 等待任意一个任务完成(或所有任务处理完毕)
tokio::select! {
Some(Ok(())) = join_set.join_next() => {
// 任务已退出,继续循环检查新就绪节点
}
Some((node_id, result)) = rx.recv() => {
// 收到工具执行结果
match result {
Ok(value) => dag.complete_node(node_id, value),
Err(e) => {
dag.nodes.get_mut(&node_id).unwrap().state =
ExecutionState::Failed(e.to_string());
// 可选:标记所有下游节点为失败
}
}
}
else => break,
}
}
Ok(DAGResult::from(dag))
}
}
3.4 信号量背压控制
并发执行不等于无限并发。在真实场景中,我们必须考虑:
- 下游数据库连接池有限(通常 20-50 连接)
- API 调用有速率限制(OpenAI 模型每分钟请求数限制)
- GPU 推理资源有限
// 多级信号量:全局并发 + 工具类型级并发
struct TieredConcurrency {
global: Arc<Semaphore>, // 全局并发上限:64
by_tool_type: HashMap<String, Arc<Semaphore>>, // 各工具类型限制
}
impl TieredConcurrency {
async fn acquire(&self, tool_name: &str) -> Result<TieredPermit, ExecutorError> {
let global_permit = self.global.clone().acquire_owned().await?;
let tool_semaphore = self.by_tool_type.get(tool_name)
.unwrap_or(&self.global);
let tool_permit = tool_semaphore.clone().acquire_owned().await?;
Ok(TieredPermit {
_global: global_permit,
_tool: tool_permit,
})
}
}
// 配置示例
let concurrency = TieredConcurrency {
global: Arc::new(Semaphore::new(64)),
by_tool_type: HashMap::from([
("database_query".into(), Arc::new(Semaphore::new(16))),
("http_api_call".into(), Arc::new(Semaphore::new(32))),
("gpu_inference".into(), Arc::new(Semaphore::new(4))),
]),
};
四、实战:构建一个支持并行工具调用的 MCP Server
4.1 MCP 协议的并行语义
Model Context Protocol (MCP) 在协议层面已经支持并行工具调用。当 Server 返回 tools/call 响应时,可以包含多个工具调用请求;Client 在收到后可以并发处理。关键在于 Server 需要在工具描述中正确标注依赖关系。
# MCP Server 定义带依赖信息的工具
TOOLS = [
{
"name": "search_web",
"description": "搜索互联网获取信息",
"inputSchema": {
"type": "object",
"properties": {
"query": {"type": "string", "description": "搜索关键词"}
},
"required": ["query"]
}
},
{
"name": "analyze_content",
"description": "分析网页内容,需要 search_web 的结果作为输入",
"inputSchema": {
"type": "object",
"properties": {
"url": {"type": "string", "description": "来自 search_web 结果的 URL"},
"content": {"type": "string", "description": "来自 search_web 结果的摘要"}
},
"required": ["url", "content"]
},
# 自定义扩展:声明依赖
"_dependencies": ["search_web"]
},
{
"name": "parallel_geocode",
"description": "地理编码(独立工具,可与其他工具并行)",
"inputSchema": {
"type": "object",
"properties": {
"address": {"type": "string"}
},
"required": ["address"]
}
}
]
4.2 基于 asyncio 的 Python 执行引擎
import asyncio
import networkx as nx
from dataclasses import dataclass, field
from typing import Any
from enum import Enum
class NodeState(Enum):
PENDING = "pending"
READY = "ready"
RUNNING = "running"
COMPLETED = "completed"
FAILED = "failed"
@dataclass
class ToolTask:
name: str
args: dict
deps: list[str] = field(default_factory=list)
state: NodeState = NodeState.PENDING
result: Any = None
class ParallelToolEngine:
def __init__(self, tool_registry: dict, max_concurrency: int = 8):
self.tool_registry = tool_registry
self.semaphore = asyncio.Semaphore(max_concurrency)
self.dag = nx.DiGraph()
def build_dag(self, tool_calls: list[dict]):
"""从工具调用列表构建 DAG,根据 args 中的引用关系确定依赖"""
# 第一遍:添加所有节点
for call in tool_calls:
task = ToolTask(name=call["name"], args=call.get("args", {}))
self.dag.add_node(call["id"], task=task)
# 第二遍:分析参数引用,建立边
ref_pattern = re.compile(r'\{\{(?P<ref_tool_id>[\w-])\}\}')
for call in tool_calls:
args_str = json.dumps(call.get("args", {}))
for match in ref_pattern.finditer(args_str):
ref_id = match.group("ref_tool_id")
if ref_id in self.dag:
self.dag.add_edge(ref_id, call["id"]) # ref_id 必须先执行
self.dag.nodes[call["id"]]["task"].deps.append(ref_id)
async def execute(self, tool_calls: list[dict]) -> dict[str, Any]:
"""执行整个 DAG,返回所有工具的结果映射"""
self.build_dag(tool_calls)
# 初始化:无依赖的节点标记为 READY
for node_id, data in self.dag.nodes(data=True):
if not data["task"].deps:
data["task"].state = NodeState.READY
running_tasks: dict[str, asyncio.Task] = {}
while len(self.get_completed_nodes()) < len(self.dag):
# 启动所有 READY 节点的执行
for node_id, data in self.dag.nodes(data=True):
if data["task"].state == NodeState.READY:
data["task"].state = NodeState.RUNNING
running_tasks[node_id] = asyncio.create_task(
self._run_tool(node_id, data["task"])
)
if not running_tasks:
break # 全部完成或有循环依赖
# 等待任意一个执行完成
done, _ = await asyncio.wait(
running_tasks.values(), return_when=asyncio.FIRST_COMPLETED
)
for completed in done:
# 找到对应的 node_id
node_id = next(k for k, v in running_tasks.items() if v is completed)
result = await completed
task = self.dag.nodes[node_id]["task"]
task.result = result
task.state = NodeState.COMPLETED if not isinstance(result, Exception) else NodeState.FAILED
running_tasks.pop(node_id)
# 检查下游节点是否变为 READY
for successor in self.dag.successors(node_id):
succ_task = self.dag.nodes[successor]["task"]
if all(
self.dag.nodes[dep]["task"].state == NodeState.COMPLETED
for dep in succ_task.deps
):
succ_task.state = NodeState.READY
return {
data["task"].result
for node_id, data in self.dag.nodes(data=True)
}
async def _run_tool(self, node_id: str, task: ToolTask):
"""带信号量背压控制的单个工具执行"""
async with self.semaphore:
try:
handler = self.tool_registry[task.name]
resolved_args = self._resolve_args(task.args)
return await handler(**resolved_args)
except Exception as e:
return e
def _resolve_args(self, args: dict) -> dict:
"""解析 {{ref_id}} 占位符,替换为已完成工具的结果"""
resolved = {}
for key, value in args.items():
if isinstance(value, str) and value.startswith("{{") and value.endswith("}}"):
ref_id = value[2:-2].strip()
ref_result = self.dag.nodes[ref_id]["task"].result
resolved[key] = ref_result.get(key, str(ref_result))
else:
resolved[key] = value
return resolved
def get_completed_nodes(self):
return [
nid for nid, data in self.dag.nodes(data=True)
if data["task"].state in (NodeState.COMPLETED, NodeState.FAILED)
]
4.3 错误处理与部分成功策略
并行执行中单个工具的失败不应直接终止整个 DAG,而是需要分级的错误处理策略:
#[derive(Debug)]
enum FailurePolicy {
/// 传播失败:下游依赖节点取消执行
Propagate,
/// 降级执行:跳过失败节点,下游使用默认值
Degraded,
/// 重试重试:在超时或限流时自动重试
Retry { max_attempts: usize, backoff_ms: u64 },
}
impl ToolDAG {
fn handle_failure(&mut self, failed_id: ToolCallId, policy: &FailurePolicy) {
match policy {
FailurePolicy::Propagate => {
// BFS 取消所有下游节点
let mut queue = vec![failed_id];
while let Some(current) = queue.pop() {
if let Some(dependents) = self.reverse_edges.get(¤t) {
for dep_id in dependents {
if !matches!(self.nodes[dep_id].state, ExecutionState::Completed(_)) {
self.nodes[dep_id].state = ExecutionState::Failed(
"dependency cancelled".into()
);
queue.push(dep_id.clone());
}
}
}
}
}
FailurePolicy::Degraded => {
// 标记失败节点,但其依赖者仍然可以尝试使用占位符执行
// 这需要下游工具的参数支持 Option<T>
}
FailurePolicy::Retry { max_attempts, backoff_ms } => {
// 在 [_run_tool] 处实现重试循环,DAG 层面记录重试次数
}
}
}
}
五、性能优化:让并行引擎真正发挥硬件能力
5.1 工具 I/O 特性的分类调度
并非所有工具都适合同等程度的并发。根据 I/O 特性,我们将工具分为三类:
| 类型 | 特征 | 并发策略 |
|---|---|---|
| CPU-bound | 本地计算、JSON 解析、代码执行 | 并发数 = CPU 核心数 |
| I/O-bound(快) | 本地文件读取、Redis 查询 (<5ms) | 高并发(100+),使用连接池 |
| I/O-bound(慢) | HTTP API、数据库查询、模型推理 | 中等并发(16-32),配合超时降级 |
# 基于工具特性的自适应并发调度
class AdaptiveConcurrencyPool:
TOOL_CATEGORIES = {
"search_web": {"category": "io_slow", "timeout_ms": 5000, "max_concurrent": 32},
"read_file": {"category": "io_fast", "timeout_ms": 200, "max_concurrent": 128},
"run_python": {"category": "cpu_bound", "timeout_ms": 30000, "max_concurrent": 8},
}
def __init__(self):
self.semaphores: dict[str, asyncio.Semaphore] = {}
for tool_name, config in self.TOOL_CATEGORIES.items():
self.semaphores[tool_name] = asyncio.Semaphore(
config["max_concurrent"]
)
async def execute_with_policy(self, tool_name: str, coro):
sem = self.semaphores.get(tool_name, self.default_sem)
async with sem:
config = self.TOOL_CATEGORIES.get(tool_name, {})
timeout_ms = config.get("timeout_ms", 10000)
return await asyncio.wait_for(coro, timeout=timeout_ms / 1000)
5.2 结构化输出加速:减少 LLM 轮次
并行工具调用的前提是有足够多的工具可以并行。一个有洞察力的优化是:通过精心设计提示词,让 LLM 一次性返回尽可能多的工具调用,而不是每轮只返回一个。
系统提示词优化策略:
1. "在 response 中,尽可能在同一轮中同时发起所有相互独立的工具调用"
2. "仅当某个工具的输出是另一个工具的必要输入时,才将它们分配在不同轮次"
3. "如果所有需要的参数已经具备,不要等待,立即并行调用"
5.3 性能基准:同步 vs DAG 并行
以下是在真实 Agent 工作负载下的基准测试数据(模拟场景:用户要求 Agent 分析 5 个 GitHub 仓库的代码质量):
工作负载:5 个并行仓库分析任务
┌──────────────────────┬──────────┬──────────┬──────────┐
│ 指标 │ 同步循环 │ 全并行 │ DAG 优化 │
├──────────────────────┼──────────┼──────────┼──────────┤
│ 总延迟 (p50) │ 18.3s │ 6.1s │ 8.4s │
│ 总延迟 (p99) │ 25.7s │ 14.2s │ 11.8s │
│ LLM 调用总 token 数 │ 287K │ 287K │ 198K │
│ 工具执行总次数 │ 25 │ 25 │ 18 │
│ 超时/失败次数 │ 0 │ 6 │ 1 │
└──────────────────────┴──────────┴──────────┴──────────┘
说明:
- "全并行" 虽然总延迟最低,但 p99 高、失败多——API 并发过高触发限流
- "DAG 优化" 通过合并 LLM 轮次减少 token 总量,同时保持合理延迟
- 实际最优方案是 DAG 分级 + 每级内并行 + 自适应限流
六、生产级架构:从原型到可靠系统
6.1 执行状态持久化
在长时间运行的 Agent 任务中,工具执行可能跨越分钟甚至小时级别。进程重启或网络中断不应丢失执行进度。
// 使用事件溯源(Event Sourcing)持久化 DAG 执行状态
#[derive(Debug, Serialize, Deserialize)]
enum ToolExecutionEvent {
ExecutionStarted { dag_id: Uuid, timestamp: DateTime<Utc> },
NodeScheduled { node_id: ToolCallId, tool_name: String },
NodeCompleted { node_id: ToolCallId, result_hash: String, duration_ms: u64 },
NodeFailed { node_id: ToolCallId, error: String, retry_count: u32 },
ExecutionFinished { dag_id: Uuid, total_duration_ms: u64 },
}
struct PersistentDAGExecutor {
event_store: Arc<dyn EventStore>,
checkpoint_interval: Duration,
}
impl PersistentDAGExecutor {
async fn execute_with_recovery(&self, dag: ToolDAG) -> Result<DAGResult, Error> {
// 检查是否有断点续传的历史记录
if let Some(checkpoint) = self.load_latest_checkpoint(&dag.id).await {
info!("从断点恢复 DAG 执行: {:?}", checkpoint);
dag.restore_from(checkpoint);
}
// 正常执行流程
let result = self.execute_dag(dag.clone()).await?;
// 事件溯源:所有状态变更写入事件存储
for event in dag.drain_events() {
self.event_store.append(event).await?;
}
Ok(result)
}
}
6.2 Observability:可观测性三大支柱
并行执行引擎的可观测性比同步循环复杂得多——你需要追踪的是一个不断变化的任务图,而不仅仅是线性调用栈。
# OpenTelemetry 集成:为每个工具调用创建 span,并通过 trace context 关联
from opentelemetry import trace
from opentelemetry.propagate import inject, extract
tracer = trace.get_tracer("agent.tool_engine")
class ObservableExecutor:
async def execute_dag(self, dag: ToolDAG):
with tracer.start_as_current_span("dag.execution") as root_span:
root_span.set_attribute("dag.node_count", len(dag.nodes))
ready_nodes = dag.find_ready_nodes()
root_span.set_attribute("dag.initial_ready_nodes", len(ready_nodes))
tasks = []
for node_id in ready_nodes:
# 为每个节点创建子 span,保持 trace 上下文
span = tracer.start_span(f"tool.{dag.nodes[node_id].name}")
ctx = trace.set_span_in_context(span)
headers = {}
inject(headers) # 注入 trace context 到工具调用参数
tasks.append(self._execute_with_span(node_id, dag, ctx))
await asyncio.gather(*tasks, return_exceptions=True)
6.3 安全防护:并行执行中的资源耗尽攻击
当 Agent 可以并行调用工具时,恶意(或被误导的)的 LLM 输出可能触发大量并发请求,形成应用层 DDoS:
class ResourceGuard:
"""防御并行工具调用中的资源耗尽"""
def __init__(self):
self.budgets: dict[str, TokenBudget] = {} # 每个会话的预算
class TokenBudget:
def __init__(self):
self.max_tool_calls_per_invocation = 16
self.max_concurrent_network_calls = 8
self.max_total_execution_time = timedelta(seconds=120)
self.network_calls_made = 0
self.start_time = time.monotonic()
def check_and_consume(self, session_id: str, tool_name: str) -> bool:
"""返回是否允许执行该工具"""
budget = self.budgets.setdefault(session_id, self.TokenBudget())
if budget.network_calls_made >= budget.max_concurrent_network_calls:
return False
if time.monotonic() - budget.start_time > budget.max_total_execution_time.total_seconds():
return False
if tool_name.startswith("http_") or tool_name.startswith("web_"):
budget.network_calls_made += 1
return True
七、总结与展望
从同步 ReAct 循环到 DAG 异步执行引擎,AI Agent 的工具调度正在经历从"单线程脚本"到"分布式任务编排"的范式迁移。核心要点回顾:
- 依赖分析是并行的前提:通过 LLM 输出中的参数引用或显式依赖声明构建工具调用 DAG。
- 分层信号量控制并发:全局、工具类型、下游服务三级信号量避免资源耗尽。
- 错误处理是并行的衍生品:并行意味着错误传播路径从线性变为图结构,需要显式的失败策略。
- 状态持久化不可或缺:长任务必须支持断点续传,事件溯源是 DAG 状态的自然建模方式。
- 可观测性是生产门槛:从线性 span 到嵌套/并行 trace,需要 OpenTelemetry 的深度集成。
未来的方向也很清晰:随着 LLM 推理速度持续提升(Apple Silicon 上的本地推理已突破 100 tok/s),工具执行将取代模型推理成为 Agent 延迟的主要瓶颈。谁能更高效地调度并行工具调用,谁就能在 Agent 工程化竞赛中占据先机。
代码仓库示例:以上所有代码片段可在 GitHub 上的实验性项目 parallel-agent-engine 中找到完整实现。

发表评论 取消回复