并发控制与速率限制:大模型 API 高吞吐实战
并发控制是大规模调用大模型 API 的必修课:并发过低浪费吞吐量,过高触发 429 限流甚至封号。核心思路是在客户端主动限速,让实际并发始终低于平台限制,而不是等触发 429 再退避。
你大概率是踩着坑找到这篇文章的:脚本跑到第 40 条时候突然一批全部失败,日志里刷屏 RateLimitError: Error code: 429 - {'error': {'message': 'Rate limit reached for requests', 'type': 'requests', 'param': None, 'code': 'rate_limit_exceeded'}},你以为是账号被封,实际上只是并发开太猛。反过来,如果你把并发压得死死的、一次只发一条等结果回来再发下一条,跑 500 条 prompt 可能要等上十几分钟——这才是真正在浪费你花钱买的配额。并发控制要解决的就是这个中间态:既不触发限流,又尽量把配额用满。
限速维度与目标
在做并发规划前先确认平台配额:
| 配额类型 | 查看方式 | 控制策略 |
|---|---|---|
| RPM(每分钟请求数) | 平台控制台 / 响应头 | 信号量 + 滑动窗口计数 |
| TPM(每分钟 token) | 平台控制台 | 估算单次 token 量后限并发 |
| 并发连接数 | 平台文档 | asyncio Semaphore |
建议将目标并发设为平台限制的 70-80%,留出余量应对突发。
这个 70-80% 不是拍脑袋的经验数字,是因为你观测到的 429 通常有滞后——平台内部的限速计数器和你本地的计数器不是同一时刻采样,网络往返也有几十到几百毫秒延迟,满打满算跑到 100% 大概率会在瞬时抖动时越界。留 20-30% 缓冲带,本质是给这个采样误差留余地。反过来并发压到 50% 以下也没必要,那是在为不存在的风险买单,白白拖慢整体吞吐。
Python asyncio + Semaphore
import os, asyncio
from openai import AsyncOpenAI
client = AsyncOpenAI(
api_key=os.environ["OPENAI_API_KEY"],
base_url=os.environ.get("OPENAI_BASE_URL", "https://api.lidayun.com/v1"),
)
MAX_CONCURRENT = 10 # 同时在途请求数
semaphore = asyncio.Semaphore(MAX_CONCURRENT)
async def chat_one(prompt: str, model: str = "gpt-4o-mini") -> str:
async with semaphore:
resp = await client.chat.completions.create(
model=model,
messages=[{"role": "user", "content": prompt}],
max_tokens=512,
)
return resp.choices[0].message.content
async def batch_chat(prompts: list[str]) -> list[str]:
tasks = [chat_one(p) for p in prompts]
return await asyncio.gather(*tasks, return_exceptions=True)
# 使用
prompts = [f"问题 {i}:大模型 API 是什么?" for i in range(50)]
results = asyncio.run(batch_chat(prompts))
for i, r in enumerate(results):
if isinstance(r, Exception):
print(f"[{i}] 失败: {r}")
else:
print(f"[{i}] {r[:30]}…")
这段代码里 asyncio.Semaphore(10) 干的事情很直白:同一时刻最多 10 个协程能拿到锁往下走,第 11 个协程会在 async with semaphore 那一行原地挂起,等前面某个请求 return 释放锁再唤醒它。这跟开 10 个线程完全是两回事——协程本身不占用系统线程,asyncio 的事件循环用单线程调度所有协程,await 到网络 I/O 时会让出控制权给别的协程用,所以哪怕开到几百个并发协程,内存和 CPU 开销也远比多线程小。但要注意一点:AsyncOpenAI 客户端底层用的是 httpx,httpx 默认连接池上限是 100 个连接(以官方文档标注的默认值为准,不同版本可能调整),如果你把 MAX_CONCURRENT 开到超过这个数,多出来的请求会卡在连接池排队而不是报错,表现就是延迟莫名其妙变高却看不到任何异常——排查这个坑的关键是去看 httpx 的连接池配置,而不是怀疑代码逻辑有问题。
RPM 滑动窗口限速
仅靠信号量控制并发数还不够,还需限制每分钟请求总数:
import time
from collections import deque
class RateLimiter:
"""滑动窗口 RPM 限速器"""
def __init__(self, rpm: int):
self.rpm = rpm
self.window = 60.0
self.timestamps: deque[float] = deque()
async def acquire(self):
while True:
now = time.monotonic()
# 清除窗口外的记录
while self.timestamps and now - self.timestamps[0] >= self.window:
self.timestamps.popleft()
if len(self.timestamps) < self.rpm:
self.timestamps.append(now)
return
# 需等待:计算最早记录离开窗口的时间
wait = self.window - (now - self.timestamps[0]) + 0.01
await asyncio.sleep(wait)
limiter = RateLimiter(rpm=60) # 目标:60 RPM
async def chat_limited(prompt: str) -> str:
await limiter.acquire()
async with semaphore:
resp = await client.chat.completions.create(
model="gpt-4o-mini",
messages=[{"role": "user", "content": prompt}],
)
return resp.choices[0].message.content
这个限速器为什么用 deque 而不是普通 list?因为要频繁在左端弹出过期的时间戳,list 的 pop(0) 是 O(n),量一大就会拖慢限速判断本身;deque.popleft() 是 O(1),这才是它在高频调用场景下的正确选择。另外注意用的是 time.monotonic() 而不是 time.time()——time.time() 会受系统时间被 NTP 校准、手动改系统时间等因素影响,可能出现”时间倒退”,导致窗口计算出诡异的负数等待时间;monotonic 保证的是单调递增,专门用来测量时间间隔。如果你把这个限速器部署到多进程场景(比如用 multiprocessing 或者多个 Python 进程跑同一个任务),会发现完全不生效——因为 timestamps 是进程内存里的普通变量,各进程互不可见,这时候就得挪到 Redis 里用 INCR 加过期时间实现,或者干脆把限速收敛到网关层统一做。
TPM 限速:请求数不是唯一瓶颈
很多人只盯着 RPM 调参,却忽略了 TPM(每分钟 token 数)经常才是真正先被打满的那道墙——尤其是长文本摘要、RAG 场景动辄几千 token 的 prompt,10 个并发请求,每个 prompt 4000 token,理论峰值就是 40000 token 塞进一分钟里,很容易先撞上 TPM 上限而不是 RPM 上限。这种场景光靠限并发数不够,得先估算每条请求大概花多少 token,再倒推能开多大并发。
import tiktoken
enc = tiktoken.encoding_for_model("gpt-4o-mini")
def estimate_tokens(text: str) -> int:
return len(enc.encode(text))
def safe_concurrency(tpm_limit: int, avg_tokens_per_request: int, avg_latency_sec: float = 3.0) -> int:
"""根据 TPM 配额和单次请求的平均耗时,反推安全并发数"""
tokens_per_second_budget = tpm_limit / 60
tokens_per_request_per_second = avg_tokens_per_request / avg_latency_sec
return max(1, int(tokens_per_second_budget / tokens_per_request_per_second * 0.75))
# 示例:TPM 上限 200000,单条 prompt+回复约 1500 token,平均耗时 3 秒
print(safe_concurrency(tpm_limit=200000, avg_tokens_per_request=1500, avg_latency_sec=3.0))
这段代码不是让你把结果当成绝对准确的答案,而是给你一个可执行的起点:先用 tiktoken 把历史请求的 token 分布跑一遍,拿到平均值和 P95,再用 P95 而不是平均值去算安全并发——这样即便偶尔来几条长 prompt,也不至于一下子把 TPM 打穿。跑起来之后盯着响应头里类似 x-ratelimit-remaining-tokens 的字段(不同平台命名可能不同,具体以官方文档为准)动态校准这个数字,比死算一次更靠谱。
Node.js:p-limit 控制并发
import OpenAI from "openai";
import pLimit from "p-limit"; // npm i p-limit
const client = new OpenAI({
apiKey: process.env.OPENAI_API_KEY,
baseURL: process.env.OPENAI_BASE_URL ?? "https://api.lidayun.com/v1",
});
const limit = pLimit(10); // 最多 10 个并发
async function chatOne(prompt) {
const resp = await client.chat.completions.create({
model: "gpt-4o-mini",
messages: [{ role: "user", content: prompt }],
});
return resp.choices[0].message.content;
}
const prompts = Array.from({ length: 50 }, (_, i) => `问题 ${i}`);
const results = await Promise.all(prompts.map((p) => limit(() => chatOne(p))));
console.log(results);
p-limit 和你手写一个计数器变量做并发控制,效果类似,但它替你处理了两件容易翻车的事:一是任务提前抛出异常时正确释放”槽位”,不会因为一个任务失败就导致后续任务永远排不上队;二是内部用 Promise 链而不是轮询定时器,没有额外的定时器开销。如果你的 Node 项目还要做更复杂的调度(比如按优先级排队、动态调整并发数),可以换成 Bottleneck 库,功能更全但学习成本也更高——单纯限并发数,p-limit 几行代码足够,没必要一上来就上重型方案。
滑动窗口 vs 令牌桶:怎么选
上面的 RateLimiter 用的是滑动窗口算法,它精确但有个特点:如果 60 秒内的配额提前用完,第 61 秒开始的请求必须硬等到窗口滑出去才能继续,请求节奏容易一顿一顿的。令牌桶(token bucket)解决的正是这个问题:桶里持续以固定速率生成令牌,请求来了先扣令牌,令牌攒够了就能连续发好几条,天然支持”平时攒着、突发时集中消耗”的场景。
| 维度 | 滑动窗口 | 令牌桶 |
|---|---|---|
| 实现复杂度 | 低,一个 deque 搞定 | 中,需要维护令牌数+补充时间戳 |
| 突发流量 | 不友好,窗口用满就死等 | 友好,允许短时突发消耗攒下的令牌 |
| 精确度 | 高,严格对齐平台的时间窗口定义 | 略有误差,取决于补充频率设置 |
| 适用场景 | 对标平台按分钟计费/限流的场景 | 内部任务调度、需要允许突发的场景 |
实操上给你个判断依据:如果你是直接对接大模型平台的 RPM/TPM 限制,滑动窗口更贴合平台的计费口径,出问题时对账也更好解释;如果是你自己系统内部给不同业务方分配调用额度,令牌桶的突发容忍度体验更好。没有绝对的”更好”,看你要贴合的是谁的规则。
批量处理最佳实践
| 技巧 | 说明 |
|---|---|
| 分批提交 | 将任务切成 batch_size=20 的批次,批次间 sleep 1s |
| 重试队列 | 失败任务放入重试队列,避免阻塞主流程 |
| 进度持久化 | 已完成结果写文件,程序崩溃可断点续跑 |
| token 估算 | 用 tiktoken 预估每条 prompt token 数,动态调整并发 |
常见问题
信号量设多少合适? 从 5 开始,观察 429 频率和延迟,逐步上调至平台并发限制的 70%。通常 10-20 对大多数付费账号足够。
asyncio.gather 里有一个任务失败会影响其他任务吗? 默认不会,gather 继续等待其他任务完成,失败的以异常对象返回(需传 return_exceptions=True)。
多进程/多机器时如何共享限速? 需借助 Redis 实现分布式令牌桶,或在网关层统一限速(推荐,如 API 网关的 rate limit 中间件)。
平台账号升级后如何动态调整并发? 将并发参数外置为环境变量(如 MAX_CONCURRENT=20),重启进程即生效,无需改代码。
为什么设了 Semaphore 还是偶尔报 429? 大概率是你只控制了并发数,没控制 RPM。比如 MAX_CONCURRENT=10,但每个请求平均只要 200ms 就返回,理论上一分钟能跑 10 / 0.2 * 60 = 3000 次请求,如果平台 RPM 限制是 500,这个并发数在响应快的场景下照样会把你打进限流——这也是本文把 Semaphore 和滑动窗口两层限速叠加使用的原因,只用其中一层在响应快的场景下是不够的。
重试时要不要跳过限速器直接重发? 不要。429 之后如果不经过限速器直接重试,等于给本来就超限的窗口又叠加了一次请求,大概率立刻再吃一次 429,进入死循环。正确做法是重试请求也要走 limiter.acquire(),并且额外读一下响应里的 Retry-After 头(如果平台返回了这个字段),按它给的秒数等待,没有这个字段再退化到指数退避策略。
更多接入基础见大模型 API 接入完全指南与接入教程专题。已触发 429 时的退避策略见 429 限流:指数退避与重试;超时连接管理见超时与连接管理。