上一辑 463 讲的是单个工具调用失败怎么办(超时、重试、降级、熔断)。但真实任务里另一类瓶颈更常见:Agent 一次要调几十个工具。比如「把这 20 个竞品官网各抓一页摘要」,串行循环每次 2 秒,总耗时 40 秒起步,用户早关页面了。这 40 秒里 CPU 大部分时间在等网络响应,纯属浪费。本篇把并发工程补齐:用 asyncio 把等待时间叠起来跑,再用并发上限和速率限制防止把下游打爆,最后处理并发特有的麻烦——部分失败怎么算。
Step 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:把同步工具包成异步函数
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())
Step 3:gather 一把并发,先跑通最快版本
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 设并发上限,别把下游打挂
免费 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」开始压测,观察错误率再调。
Step 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 个白跑
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 条。
Step 7:串起来——可复用的批量调用包装器
把上面所有机制收进一个批量执行器,带超时兜底:
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 的工具层就同时具备了恢复力与吞吐力。压测建议拿自己最常用的那个工具跑三组数字(串行、无上限并发、限流并发),贴在团队文档里——下次有人想全速并发时,有数字可摆。