Agent 不只是顺序调用工具,还可以判断依赖关系,把互不依赖的查询同时执行。
执行流程有一个比较明显的问题:慢。
假设排查订单服务故障时,需要查询三类信息:
- 最近30分钟的错误日志;
- CPU、内存和接口延迟;
- 最近一次发布记录。
如果按照顺序执行:
查询日志:3秒
查询监控:4秒
查询发布记录:2秒
总耗时:约9秒
但这三个查询互不依赖,没有必要一个接一个等待。并行执行后,理论耗时会接近最慢的那个工具,也就是4秒左右。
┌─ 查询日志 ───────┐
用户问题 ─┼─ 查询监控 ───────┼─ 汇总结果
└─ 查询发布记录 ───┘
听起来只是一个普通的并发优化,但放到 Agent 中并没有那么简单。程序需要判断哪些工具可以同时执行,哪些工具必须等待前一步结果。
先定义几个异步工具
下面用 Python 的 asyncio 模拟三个远程查询接口:
import asyncio
from datetime import datetime
async def query_service_logs(
service_name: str,
minutes: int
) -> dict:
await asyncio.sleep(3)
return {
"tool": "query_service_logs",
"service": service_name,
"minutes": minutes,
"error_count": 236,
"main_error": "Database connection pool exhausted"
}
async def query_service_metrics(
service_name: str,
minutes: int
) -> dict:
await asyncio.sleep(4)
return {
"tool": "query_service_metrics",
"service": service_name,
"cpu": "68%",
"memory": "74%",
"average_latency": "4.8s",
"database_pool_usage": "100%"
}
async def query_deployment_history(
service_name: str
) -> dict:
await asyncio.sleep(2)
return {
"tool": "query_deployment_history",
"service": service_name,
"last_deployment": "2026-07-24 09:30:00",
"version": "order-service:2.7.1",
"status": "success"
}
真实项目中,asyncio.sleep() 会被替换成数据库、日志平台或者 HTTP 接口调用。
例如使用异步 HTTP 客户端:
import httpx
async def query_service_logs(
service_name: str,
minutes: int
) -> dict:
async with httpx.AsyncClient(timeout=10) as client:
response = await client.get(
"https://log-api.example.com/search",
params={
"service": service_name,
"minutes": minutes
}
)
response.raise_for_status()
return response.json()
这里最好使用真正的异步客户端。如果在异步函数中调用同步的 requests.get(),仍然会阻塞整个事件循环,并不能获得理想的并发效果。
使用 asyncio.gather() 并行执行
最简单的并行方式是把多个协程交给 asyncio.gather():
async def collect_service_context(service_name: str):
started_at = datetime.now()
logs, metrics, deployment = await asyncio.gather(
query_service_logs(service_name, 30),
query_service_metrics(service_name, 30),
query_deployment_history(service_name)
)
elapsed = (datetime.now() - started_at).total_seconds()
return {
"logs": logs,
"metrics": metrics,
"deployment": deployment,
"elapsed_seconds": elapsed
}
async def main():
result = await collect_service_context("order-service")
print(result)
asyncio.run(main())
三个工具会同时启动,不需要等待前一个执行完成。
但这还只是写死的并行调用。真正的 Agent 无法提前确定模型会选择哪些工具,所以还要设计一个通用的工具执行器。
实现一个通用的并行工具执行器
先建立工具注册表:
TOOL_REGISTRY = {
"query_service_logs": query_service_logs,
"query_service_metrics": query_service_metrics,
"query_deployment_history": query_deployment_history
}
假设模型一次返回了多个工具调用:
tool_calls = [
{
"id": "call_001",
"name": "query_service_logs",
"arguments": {
"service_name": "order-service",
"minutes": 30
}
},
{
"id": "call_002",
"name": "query_service_metrics",
"arguments": {
"service_name": "order-service",
"minutes": 30
}
},
{
"id": "call_003",
"name": "query_deployment_history",
"arguments": {
"service_name": "order-service"
}
}
]
先封装单个工具的执行逻辑:
async def execute_one_tool(tool_call: dict) -> dict:
call_id = tool_call["id"]
tool_name = tool_call["name"]
arguments = tool_call["arguments"]
tool = TOOL_REGISTRY.get(tool_name)
if tool is None:
return {
"call_id": call_id,
"tool": tool_name,
"success": False,
"error": "工具不存在"
}
try:
result = await tool(**arguments)
return {
"call_id": call_id,
"tool": tool_name,
"success": True,
"result": result
}
except Exception as error:
return {
"call_id": call_id,
"tool": tool_name,
"success": False,
"error": str(error)
}
然后统一并行调度:
async def execute_tools_in_parallel(
tool_calls: list[dict]
) -> list[dict]:
tasks = [
asyncio.create_task(execute_one_tool(tool_call))
for tool_call in tool_calls
]
return await asyncio.gather(*tasks)
调用方式如下:
async def main():
results = await execute_tools_in_parallel(tool_calls)
for result in results:
print(result)
asyncio.run(main())
这样,无论模型一次返回两个工具还是十个工具,都可以使用相同的方式执行。
一个工具失败,不应该拖垮全部任务
asyncio.gather() 默认情况下,只要某个协程抛出异常,异常就会向外传播。
在 Agent 场景中,这通常不是我们想要的结果。
例如监控平台暂时不可用,但日志和发布记录查询成功了。此时应该把已有结果交给模型,而不是让整个任务直接失败。
可以使用:
results = await asyncio.gather(
*tasks,
return_exceptions=True
)
然后逐个处理结果:
async def execute_tools_in_parallel(
tool_calls: list[dict]
) -> list[dict]:
tasks = [
asyncio.create_task(execute_one_tool(tool_call))
for tool_call in tool_calls
]
raw_results = await asyncio.gather(
*tasks,
return_exceptions=True
)
results = []
for tool_call, raw_result in zip(tool_calls, raw_results):
if isinstance(raw_result, Exception):
results.append({
"call_id": tool_call["id"],
"tool": tool_call["name"],
"success": False,
"error": str(raw_result)
})
else:
results.append(raw_result)
return results
不过,如果 execute_one_tool() 本身已经捕获了所有异常,return_exceptions=True 更像是一层兜底,用来防止执行器内部出现没有预料到的错误。
给工具增加超时,避免一直等下去
Agent 经常需要调用外部系统,而外部系统不一定稳定。
如果一个接口迟迟不返回,整个任务可能一直处于等待状态。因此,每次工具调用都应该设置超时时间:
async def execute_one_tool(
tool_call: dict,
timeout_seconds: int = 10
) -> dict:
call_id = tool_call["id"]
tool_name = tool_call["name"]
arguments = tool_call["arguments"]
tool = TOOL_REGISTRY.get(tool_name)
if tool is None:
return {
"call_id": call_id,
"tool": tool_name,
"success": False,
"error": "工具不存在"
}
try:
result = await asyncio.wait_for(
tool(**arguments),
timeout=timeout_seconds
)
return {
"call_id": call_id,
"tool": tool_name,
"success": True,
"result": result
}
except asyncio.TimeoutError:
return {
"call_id": call_id,
"tool": tool_name,
"success": False,
"error": f"工具执行超过{timeout_seconds}秒"
}
except Exception as error:
return {
"call_id": call_id,
"tool": tool_name,
"success": False,
"error": str(error)
}
工具超时后,应该把失败信息交回模型:
{
"tool": "query_service_metrics",
"success": false,
"error": "工具执行超过10秒"
}
模型可以根据剩余信息继续分析,也可以告诉用户当前结论缺少监控数据支持。
不要把失败结果伪装成空数据。否则模型可能会把“查询失败”理解成“没有异常”。
并行不等于无限并发
假设模型一次生成100个查询请求,如果程序全部同时执行,很容易把下游系统打挂。
因此还需要使用信号量限制并发数量:
MAX_CONCURRENCY = 5
tool_semaphore = asyncio.Semaphore(MAX_CONCURRENCY)
async def execute_one_tool(
tool_call: dict,
timeout_seconds: int = 10
) -> dict:
async with tool_semaphore:
call_id = tool_call["id"]
tool_name = tool_call["name"]
arguments = tool_call["arguments"]
tool = TOOL_REGISTRY.get(tool_name)
if tool is None:
return {
"call_id": call_id,
"tool": tool_name,
"success": False,
"error": "工具不存在"
}
try:
result = await asyncio.wait_for(
tool(**arguments),
timeout=timeout_seconds
)
return {
"call_id": call_id,
"tool": tool_name,
"success": True,
"result": result
}
except asyncio.TimeoutError:
return {
"call_id": call_id,
"tool": tool_name,
"success": False,
"error": f"工具执行超过{timeout_seconds}秒"
}
except Exception as error:
return {
"call_id": call_id,
"tool": tool_name,
"success": False,
"error": str(error)
}
这样即使模型生成了大量工具调用,系统也只会同时执行5个。
具体并发数量不能拍脑袋决定,需要结合:
- 下游接口的限流规则;
- 数据库连接池大小;
- 单次请求平均耗时;
- Agent 同时在线的任务数量;
- 模型生成工具调用的上限。
单个任务并发5次看起来不多,但如果同时有100个Agent任务运行,实际可能产生500个请求。
哪些工具可以并行,哪些不能?
这是并行 Agent 中最容易踩坑的地方。
下面三个工具互不依赖,可以同时执行:
查询日志
查询监控
查询发布记录
但下面这组操作不能直接并行:
查询日志
↓
判断故障等级
↓
创建故障工单
↓
发送工单通知
创建工单需要使用前面的分析结果,发送通知又需要工单编号。它们之间存在明确的依赖关系。
可以把整个任务看成一张有向无环图,也就是 DAG:
查询日志 ─┐
查询监控 ─┼─> 分析故障 ─> 创建工单 ─> 发送通知
查询发布 ─┘
第一层的三个查询可以并行,后面的步骤需要按依赖顺序执行。
代码可以写成:
async def handle_incident(service_name: str):
logs_task = query_service_logs(service_name, 30)
metrics_task = query_service_metrics(service_name, 30)
deployment_task = query_deployment_history(service_name)
logs, metrics, deployment = await asyncio.gather(
logs_task,
metrics_task,
deployment_task
)
analysis = await analyze_incident(
logs=logs,
metrics=metrics,
deployment=deployment
)
if not analysis["should_create_ticket"]:
return {
"analysis": analysis,
"ticket": None
}
ticket = await create_incident_ticket(
title=analysis["title"],
level=analysis["level"],
description=analysis["description"]
)
notification = await send_ticket_notification(
ticket_id=ticket["ticket_id"],
level=analysis["level"]
)
return {
"analysis": analysis,
"ticket": ticket,
"notification": notification
}
这里的原则很简单:
没有数据依赖的工具可以并行,有数据依赖的步骤必须等待。
真正困难的不是写出 asyncio.gather(),而是正确识别依赖关系。
不要并行执行多个高风险写操作
查询类工具通常适合并行:
读取日志、查询指标、搜索文档、读取发布记录
写操作则要谨慎:
创建工单、修改配置、重启服务、发送消息、删除数据
假设模型同时生成:
重启服务
修改数据库连接池配置
回滚应用版本
这三个动作如果并行执行,可能互相影响,最后甚至无法判断是哪一个操作解决了问题。
我在实现时会给工具增加元数据:
TOOL_METADATA = {
"query_service_logs": {
"read_only": True,
"parallel_safe": True
},
"query_service_metrics": {
"read_only": True,
"parallel_safe": True
},
"create_incident_ticket": {
"read_only": False,
"parallel_safe": False
},
"restart_production_service": {
"read_only": False,
"parallel_safe": False
}
}
调度前先分组:
def group_tool_calls(tool_calls: list[dict]):
parallel_calls = []
sequential_calls = []
for tool_call in tool_calls:
metadata = TOOL_METADATA.get(
tool_call["name"],
{
"read_only": False,
"parallel_safe": False
}
)
if metadata["parallel_safe"]:
parallel_calls.append(tool_call)
else:
sequential_calls.append(tool_call)
return parallel_calls, sequential_calls
安全工具并行执行,高风险工具按顺序处理:
async def execute_tool_batch(tool_calls: list[dict]):
parallel_calls, sequential_calls = group_tool_calls(
tool_calls
)
results = []
if parallel_calls:
parallel_results = await execute_tools_in_parallel(
parallel_calls
)
results.extend(parallel_results)
for tool_call in sequential_calls:
result = await execute_one_tool(tool_call)
results.append(result)
return results
实际生产环境中,非只读工具还应该经过权限检查或人工确认,不能仅仅因为模型调用了它就直接执行。
重试不能随便加
接口偶尔超时,可以增加重试,但不是所有操作都适合重试。
查询日志失败后重新查询,一般问题不大。创建工单失败后直接重试,则可能生成两张相同工单。
可以先为工具标记是否允许重试:
TOOL_METADATA = {
"query_service_logs": {
"parallel_safe": True,
"retryable": True
},
"create_incident_ticket": {
"parallel_safe": False,
"retryable": False
}
}
对查询类工具实现简单重试:
async def execute_with_retry(
tool,
arguments: dict,
max_attempts: int = 3
):
last_error = None
for attempt in range(1, max_attempts + 1):
try:
return await tool(**arguments)
except Exception as error:
last_error = error
if attempt < max_attempts:
await asyncio.sleep(2 ** (attempt - 1))
raise last_error
等待时间依次为1秒、2秒,避免接口故障时立即连续请求。
如果确实要重试写操作,就需要加入幂等键:
import hashlib
import json
def create_idempotency_key(
task_id: str,
tool_name: str,
arguments: dict
) -> str:
raw = json.dumps(
{
"task_id": task_id,
"tool": tool_name,
"arguments": arguments
},
sort_keys=True,
ensure_ascii=False
)
return hashlib.sha256(raw.encode("utf-8")).hexdigest()
调用工单接口时携带该键:
idempotency_key = create_idempotency_key(
task_id="task_20260724_001",
tool_name="create_incident_ticket",
arguments=arguments
)
服务端发现相同的幂等键后,应该返回第一次创建的结果,而不是再次创建数据。
把并行执行接入 Agent Loop
最后,把工具调度器接入前面的智能体循环:
async def run_agent(user_input: str):
messages = [
{
"role": "system",
"content": """
你是一名技术故障处理助手。
处理任务时遵守以下要求:
1. 互不依赖的只读查询可以同时调用;
2. 写操作必须等待必要的查询和分析结果;
3. 工具调用失败时,不能把失败解释为没有异常;
4. 高风险操作必须等待用户确认;
5. 最终结论必须说明依据。
"""
},
{
"role": "user",
"content": user_input
}
]
for step in range(8):
response = await llm.chat(
messages=messages,
tools=tools
)
tool_calls = response.get("tool_calls", [])
if not tool_calls:
return response["content"]
messages.append(response)
tool_results = await execute_tool_batch(tool_calls)
for result in tool_results:
messages.append({
"role": "tool",
"tool_call_id": result["call_id"],
"name": result["tool"],
"content": json.dumps(
result,
ensure_ascii=False
)
})
return "任务执行步骤超过限制,已停止继续调用工具。"
这里有个细节不能忽略:每个工具结果都要带回对应的 tool_call_id。
并行执行时,工具完成顺序未必和调用顺序一致。调用编号可以帮助模型正确对应“哪个结果属于哪个请求”。
并行工具调用到底能带来什么?
完成并行调度后,最直接的变化是速度。
原来一个故障排查任务需要依次查询日志、监控和发布记录,可能十几秒才能拿到完整结果。改成并行后,等待时间通常取决于最慢的那个查询。
但它的价值不只是“快几秒”。
并行工具调用意味着Agent可以像一个有经验的工程师一样,同时从多个方向收集证据:
日志告诉它发生了什么;
监控告诉它影响有多大;
发布记录告诉它最近改了什么;
知识库告诉它以前如何处理。
等这些结果回来后,再统一分析,而不是查到一条信息就急着下结论。
如果想直接体验这种“让 AI 调用工具完成工作”的方式,也可以试试 SnailAgent。对普通用户来说,不需要先理解并发调度、工具注册和执行循环;对开发者来说,则可以从实际使用中观察:AI 是如何拆解任务、调用工具并根据结果继续行动的。
当然,并行并不是越多越好。
一个可靠的Agent调度器至少要处理好四件事:
依赖关系、并发上限、失败隔离、操作安全
asyncio.gather() 只解决了“同时执行”的问题。真正决定系统能不能上线的,是后面这些看起来不够炫、却不能省略的工程细节。
更多推荐


所有评论(0)