进阶 📋 7 个步骤 第 465 / 467 篇

Agent 并发工具调用工程:Semaphore 限流、结果聚合与部分失败处理

把串行 20 次搜索从 40 秒压到 4 秒:asyncio.gather 并发、Semaphore 设并发上限、令牌桶控速率、聚合时区分成败。与容错三板斧(463)互补,一篇解决 Agent 批量调工具的吞吐问题。

2026.09.23· 16 分钟阅读· 约 2101 字· 🐍 Python / ⚡ asyncio

上一辑 463 讲的是单个工具调用失败怎么办(超时、重试、降级、熔断)。但真实任务里另一类瓶颈更常见:Agent 一次要调几十个工具。比如「把这 20 个竞品官网各抓一页摘要」,串行循环每次 2 秒,总耗时 40 秒起步,用户早关页面了。这 40 秒里 CPU 大部分时间在等网络响应,纯属浪费。本篇把并发工程补齐:用 asyncio 把等待时间叠起来跑,再用并发上限和速率限制防止把下游打爆,最后处理并发特有的麻烦——部分失败怎么算。

💡 前置知识只需一条:Python 3.11+ 自带 asyncio。如果你还没有把工具写成 async 函数的习惯,本篇 Step 2 会从同步函数的包装讲起,照抄即可跑通。

Step 1:先量化问题——串行到底慢在哪

1 写一个最朴素的串行版本当对照组

用一个模拟搜索工具(内含 sleep 模拟网络往返)跑 8 次调用,记录总耗时:

import time

def search(query: str) -> str:
    time.sleep(0.5)          # 模拟一次 500ms 的网络请求
    return f"【{query}】的搜索结果摘要"

queries = [f"竞品{i}官网" for i in range(8)]

t0 = time.perf_counter()
results = [search(q) for q in queries]
print(f"串行耗时: {time.perf_counter() - t0:.2f}s")
# 典型输出: 串行耗时: 4.03s

8 次 0.5 秒的调用串行跑了 4 秒。换成真实的网页抓取或第三方 API,单次延迟经常在 1 到 3 秒,20 个目标的串行版本会慢到不可用。

Step 2:把同步工具包成异步函数

2 同步函数用 to_thread 塞进线程池

asyncio 的前提是「可等待对象」。已经写成 async 的 HTTP 客户端(如 httpx.AsyncClient)直接用;同步函数则用 asyncio.to_thread 包一层:

import asyncio

async def asearch(query: str) -> str:
    # to_thread 把阻塞调用丢进线程池,不卡事件循环
    return await asyncio.to_thread(search, query)

# 冒烟测试
async def main():
    r = await asearch("测试")
    print(r)

asyncio.run(main())
💡 什么情况能不用 to_thread?如果你的工具本身是异步客户端(httpx.AsyncClient、aiohttp、asyncpg),直接 await 它即可,性能更好。to_thread 是给存量同步代码的通用解法,一次改动全库受益。

Step 3:gather 一把并发,先跑通最快版本

3 一行 gather,40 秒变 2 秒

asyncio.gather 接收一组协程,同时调度它们,全部完成后按原顺序返回结果列表:

async def run_all(queries):
    t0 = time.perf_counter()
    results = await asyncio.gather(*(asearch(q) for q in queries))
    print(f"并发耗时: {time.perf_counter() - t0:.2f}s")
    return results

asyncio.run(run_all(queries))
# 典型输出: 并发耗时: 0.51s —— 8 个请求几乎同时发出,总耗时约等于最慢的那一个

耗时从 4 秒降到 0.5 秒。但先别高兴——这个版本有两个隐患:任何一次抛异常会炸掉整个 gather;8 个、80 个、800 个并发对下游是三种完全不同的冲击。接下来逐个解决。

asyncio 是单线程并发,只加速「等待 IO」的任务。如果工具内部是纯 CPU 计算(大文件解析、图像处理),gather 不会变快,反而要上 multiprocessing 或进程池。判断标准很简单:耗时主要花在等网络、等磁盘,用并发;花在算,换进程。

Step 4:Semaphore 设并发上限,别把下游打挂

4 全速并发是给别人添堵,给自己埋雷

免费 API 常限每秒几次请求,自建服务也有连接池上限。不加控制的全量并发轻则被限流拒绝,重则触发封禁。用 asyncio.Semaphore 把「同时在飞」的请求数钉在固定值:

import asyncio, random

SEM = asyncio.Semaphore(4)   # 最多 4 个请求同时在飞

async def asearch_limited(query: str) -> str:
    async with SEM:                      # 拿到令牌才执行
        return await asyncio.to_thread(search, query)

async def main():
    t0 = time.perf_counter()
    results = await asyncio.gather(*(asearch_limited(q) for q in queries))
    print(f"受限并发耗时: {time.perf_counter() - t0:.2f}s")
    print(f"完成数: {len(results)}")

asyncio.run(main())
# 典型输出: 受限并发耗时: 1.03s —— 8 个请求分两批跑完,速度仍是串行的 4 倍

并发上限 4 时,8 个请求分两波,总耗时约两倍单次延迟——依然远快于串行,且下游压力可控。上限取多少没有公式,建议从「下游限流阈值除以 2」开始压测,观察错误率再调。

💡 Semaphore 的取值可以按工具分别配置:搜网页开 8,调数据库开 3,调贵重模型开 1。把上限写进工具注册表(461 的 TOOLS 结构加一个 max_concurrency 字段),而不是散落在业务代码里。

Step 5:令牌桶控速率,压住「每秒请求数」这条线

5 并发上限管「同时多少」,速率限制管「总共多快」

Semaphore 限制同时 inflight 数,但请求极快时每秒仍然可能打出大量请求。很多 API 按每秒请求数(RPS)限流,这时需要令牌桶:

import asyncio, time

class TokenBucket:
    def __init__(self, rate: float, capacity: int):
        self.rate = rate            # 每秒补充令牌数
        self.capacity = capacity    # 桶容量
        self.tokens = float(capacity)
        self.updated = time.monotonic()
        self._lock = asyncio.Lock()

    async def acquire(self):
        async with self._lock:
            while True:
                now = time.monotonic()
                self.tokens = min(self.capacity,
                                  self.tokens + (now - self.updated) * self.rate)
                self.updated = now
                if self.tokens >= 1:
                    self.tokens -= 1
                    return
                # 令牌不够,睡到攒够一个再试
                await asyncio.sleep((1 - self.tokens) / self.rate)

BUCKET = TokenBucket(rate=5, capacity=5)   # 每秒最多 5 个请求

async def asearch_throttled(query: str) -> str:
    async with SEM:
        await BUCKET.acquire()
        return await asyncio.to_thread(search, query)

令牌桶是进程内的,多实例部署时各进程各有一个桶,全局限流要靠 Redis + 分布式限流器或网关层(Nginx limit_req、云厂商限流配置)兜底。别把进程内限流当成集群限流。

Step 6:部分失败处理——80 个里挂 3 个,别让 77 个白跑

6 return_exceptions 把失败降级成「一条结果」

gather 默认行为是快失败:一个协程抛异常,整体取消。批量场景要的是 return_exceptions=True,让异常作为元素返回,再聚合时分类:

async def run_batch(queries):
    tasks = [asearch_limited(q) for q in queries]
    raw = await asyncio.gather(*tasks, return_exceptions=True)
    ok, failed = [], []
    for q, r in zip(queries, raw):
        if isinstance(r, Exception):
            failed.append({"query": q, "error": str(r)})
        else:
            ok.append({"query": q, "result": r})
    return ok, failed

async def main():
    ok, failed = await run_batch(queries)
    print(f"成功 {len(ok)} 条, 失败 {len(failed)} 条")
    # 把失败项交给上层决策:重试、换工具,或如实告知用户缺失部分

asyncio.run(main())

聚合后把 ok 与 failed 都交给 Agent 主循环:失败项可以带着错误信息重跑一轮(注意用 465 的并发上限,别重试也全速),或者像 463 的降级策略那样换备用工具。原则是「部分成功也是可交付的成果」,别追求全对而丢掉已有的 77 条。

💡 对批量结果做日志时,建议记录 (query, 耗时, 成败) 三元组。跑一段时间后你会得到每个工具的真实延迟分布,Semaphore 上限和超时参数就从拍脑袋变成有据可调。

Step 7:串起来——可复用的批量调用包装器

7 一个函数收口,Agent 主循环只管往里塞任务

把上面所有机制收进一个批量执行器,带超时兜底:

async def batch_call(fn, args_list, sem_size=4, per_timeout=10):
    sem = asyncio.Semaphore(sem_size)

    async def one(args):
        async with sem:
            try:
                return await asyncio.wait_for(fn(*args), timeout=per_timeout)
            except Exception as e:
                return e                       # 异常当结果返回,不打断整体

    raw = await asyncio.gather(*(one(a) for a in args_list))
    return [
        {"args": a, "ok": not isinstance(r, Exception), "value": r}
        for a, r in zip(args_list, raw)
    ]

# 用法:把 20 个抓取任务一次交给它
tasks = [(q,) for q in queries]
out = asyncio.run(batch_call(asearch, tasks, sem_size=4, per_timeout=10))
print(sum(1 for x in out if x["ok"]), "/", len(out), "succeeded")

并发 + 超时组合有一个隐蔽坑:wait_for 超时后协程被取消,但如果工具内部持有未释放的资源(数据库连接、文件句柄),取消不等于释放。写异步工具时把资源获取放在 async with 里,让取消路径也能走清理逻辑。

收个尾:463 解决「单次调用挂了怎么办」,本篇解决「大量调用怎么跑得快且不闯祸」。两者组合后,Agent 的工具层就同时具备了恢复力与吞吐力。压测建议拿自己最常用的那个工具跑三组数字(串行、无上限并发、限流并发),贴在团队文档里——下次有人想全速并发时,有数字可摆。

← 返回教程中心