执行流程有一个比较明显的问题:慢。

假设排查订单服务故障时,需要查询三类信息:

  • 最近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() 只解决了“同时执行”的问题。真正决定系统能不能上线的,是后面这些看起来不够炫、却不能省略的工程细节。

蜗牛智能体以大模型为核心的办公智能体:自然语言驱动、零技术门槛、任务自动执行并交付结果。多模式(问答/执行/配置/自主)、多智能体协同、技能库与 IM 通道。数据留在本机,可选云端同步。https://snailagent.com/

Logo

中国智能体开发者社区,聚焦智能体与大模型开发,提供前沿资讯、实用工具链、开源项目及行业案例。通过技术沙龙、开发者大赛等活动,促进经验交流与协作,助力开发者快速构建创新智能应用。

更多推荐