在处理单次 AI 请求时,同步代码简单好用。但如果要处理批量数据或多轮并发对话,继续使用同步请求会让程序大量时间消耗在“干等”上。本文介绍一种适合个人项目和工具站的异步改造思路。
为什么要做异步改造?
想象一个场景:你需要用 AI 给 100 篇文章生成摘要。
如果用同步代码(requests 或 openai 默认的同步客户端),流程是这样的:
发请求 1 → 等 2 秒 → 收到结果 1 → 发请求 2 → 等 2 秒 → 收到结果 2
100 篇文章至少需要等 200 秒。在这 200 秒里,你的 CPU 几乎什么都没做,全在等待网络响应。
而异步代码的思路是:
发出请求 1 → 不等它回来,直接发请求 2 → 发请求 3 …
当某个请求的结果回来了,再去处理它。
这样,你的总耗时将取决于最慢的那几个请求,而不是所有请求耗时的总和。
一、同步调用的瓶颈在哪里?
先看一段典型的同步调用代码:
import time
from openai import OpenAI
client = OpenAI(
api_key="your-api-key",
base_url="https://your-api-domain.com/v1"
)
def process_items(items: list[str]):
results = []
start = time.perf_counter()
for item in items:
# 这里会阻塞,直到收到响应
resp = client.chat.completions.create(
model="your-model-name",
messages=[{"role": "user", "content": item}]
)
results.append(resp.choices[0].message.content)
print(f"总耗时: {time.perf_counter() – start:.2f}s")
return results
如果 items 里有 10 个元素,每个接口响应 2 秒,整个函数需要执行约 20 秒。
二、使用 AsyncOpenAI 替换同步客户端
要实现异步请求,首先需要使用支持异步的网络客户端。OpenAI 的 Python SDK 原生提供了 AsyncOpenAI。
1. 初始化异步客户端
import asyncio
from openai import AsyncOpenAI
aclient = AsyncOpenAI(
api_key="your-api-key",
base_url="https://your-api-domain.com/v1"
)
2. 定义异步请求函数
把普通的 def 改成 async def,并在网络请求前加上 await:
async def ask_llm_async(prompt: str) –> str:
response = await aclient.chat.completions.create(
model="your-model-name",
messages=[{"role": "user", "content": prompt}]
)
return response.choices[0].message.content
这里的 await 是关键:它告诉 Python,在等待接口响应的时候,可以去执行其他任务。
三、使用 asyncio.gather 并发执行
把单个请求包装成异步函数后,可以使用 asyncio.gather 同时发起多个请求:
async def process_items_async(items: list[str]):
start = time.perf_counter()
# 创建所有任务
tasks = [ask_llm_async(item) for item in items]
# 并发执行,等待所有任务完成
results = await asyncio.gather(*tasks)
print(f"异步总耗时: {time.perf_counter() – start:.2f}s")
return results
# 运行异步函数
if __name__ == "__main__":
items = ["介绍 Python", "介绍 Java", "介绍 Rust"]
asyncio.run(process_items_async(items))
现在,即使有 10 个请求,它们几乎也是同时发出的。总耗时可能只有 2 到 3 秒。
四、直接并发太猛了怎么办?加入并发限制
上一节的代码有一个隐患:如果 items 里有 1000 个元素,asyncio.gather 会瞬间发起 1000 个连接。
这可能导致:
- 你的电脑网络句柄耗尽
- 触发目标 API 的高频限流(HTTP 429 Rate Limit)
- 服务端直接报错
所以,并发一定要加以限制。可以使用 asyncio.Semaphore。
结合 Semaphore 控制并发
# 最多允许 5 个任务同时进行
concurrency_limit = asyncio.Semaphore(5)
async def safe_ask_llm(prompt: str) –> str:
# 获取信号量,拿到名额的任务才能继续往下走
async with concurrency_limit:
return await ask_llm_async(prompt)
async def safe_process_batch(items: list[str]):
tasks = [safe_ask_llm(item) for item in items]
return await asyncio.gather(*tasks)
这样既提升了效率,又保证了请求不会像洪水一样把系统冲垮。
五、在异步循环中增加容错与延迟
如果在并发请求中遇到了 RateLimitError,可以加入简单的等待重试机制,或者设置任务间的启动间隔。
async def safe_ask_with_retry(prompt: str, retries: int = 3) –> str:
for attempt in range(1, retries + 1):
try:
async with concurrency_limit:
return await ask_llm_async(prompt)
except Exception as e:
if attempt < retries:
# 遇到错误时,让当前协程休息一会儿
await asyncio.sleep(2 * attempt)
else:
return f"Error: {e}"
注意:这里用的是 await asyncio.sleep() 而不是 time.sleep()。在异步代码中,千万不要使用 time.sleep(),否则会把所有并发任务一起卡死。
六、处理部分任务失败:不要一错全挂
默认情况下,如果 asyncio.gather 里的某一个任务抛出了异常,整个 gather 都会报错退出,导致其他已经成功的结果也拿不到。
如果要避免这种情况,可以使用 return_exceptions=True:
async def process_and_keep_errors(items: list[str]):
tasks = [safe_ask_with_retry(item) for item in items]
# 某个任务报错,会把 Exception 对象放在 results 列表中对应位置
results = await asyncio.gather(*tasks, return_exceptions=True)
for i, res in enumerate(results):
if isinstance(res, Exception):
print(f"任务 {i} 失败: {res}")
else:
print(f"任务 {i} 成功")
return results
七、哪些场景特别需要异步改造?
并不是所有代码都需要加 async。如果你的工具每次只和用户一问一答,同步代码完全足够。
真正需要异步的场景:
- AI 批量处理脚本:比如批量翻译、批量提取文档摘要。
- Agent 工作流:一个任务需要并发分发给多个“专家模型”评估,然后再汇总。
- 高并发后端服务:如果你用 FastAPI 搭建自己的 AI 工具站,底层调用必须是异步的,否则几个用户提问就会把服务器线程占满。
八、结语
异步改造的核心在于:不让宝贵的 CPU 时间浪费在等待网络响应上。
完整的改造思路可以总结为三步:
一旦跑通了这个流程,你会发现处理几百条 AI 任务的时间,可能比之前快了十几倍。
免责声明
本文内容仅用于技术交流与经验分享,具体实现请结合项目实际并发需求和上游接口限制调整。



