一千个请求一起丢出去,十分钟后日志开始刷同一行:openai.RateLimitError: Error code: 429,紧跟的服务端原文是 Error code: 429 - {'error': {'message': 'Too Many Requests'}}。有人会怀疑令牌失效、余额不足,方向都不对。上游对每个账号能同时处理多少条请求、每分钟能接收多少条请求都设了闸门,超出的部分原样退回,能调的只有你自己的发送节奏。
并发数指同一时刻服务端正在处理的请求条数,RPM 指每分钟接收的请求条数。两者互相独立,谁先被顶到谁先报错。你可能并发只有 5,看起来很老实,但每个请求 30 毫秒就返回,一分钟仍然能打出去几千条,撞的是 RPM 那道门;单条要跑 40 秒的推理模型,并发开 20 就撞并发那道门。
429 的意思是「你发太快了,我在,只是现在不接」,请求本身没问题,等一会儿重发就行。把 429 和上游故障混在一起无脑重试,是脚本越跑越糟的常见起点:429 之后马上重发,等于往已经堵住的管道里继续塞东西,成功率只会更低。
没有节奏控制的脚本在 429 上有个典型特征:错误行密密麻麻,时间戳间隔极短,重试次数爆炸。你写的是「失败重试三次」,实际效果是同一秒里几百条请求在互相踩。判断方法很简单:失败条数除以总条数超过一成,就先降并发,而不是加机器。
asyncio.gather 一次造一千个协程很多人以为 await asyncio.gather(*[task(i) for i in range(1000)]) 是「排队执行」。真实过程是:列表推导式先把 1000 个协程对象全部创建好,gather 一调用,事件循环立刻把它们整体推进到发请求这一步,结果是 1000 个请求几乎同时建连、几乎同时到达上游。这不是批量处理,是给自己上游做压力测试。
上游在极短时间内收到远超限额的流量,会把超出的部分全部 429 掉。客户端看到失败就重试,重试落在同一批积压请求后面,超时从 60 秒抬到 90 秒,连接池被打满,新请求排不上队。一个本来 40 分钟能干完的任务,三小时没收尾,钱还烧掉一倍。
在途请求指已经发出去、还没拿到响应的那些。把它摁在上游能吃下的范围内,剩下的任务在本地排队等着,改动通常只要十几行代码。
asyncio.Semaphore 摁住在途请求数把 Semaphore(8) 想成八把钥匙挂在门口。每个协程发请求前先拿一把,拿到才发,拿到响应把钥匙还回去。第九个协程到门口发现钥匙发完了,就地等着,直到有人还回来。上游永远看不到超过 8 条在途请求,你本地的任务照样全部跑完。
# sem_batch.py —— 信号量控制并发,异常自己兜住,成功失败分开计数
import asyncio
import os
from openai import AsyncOpenAI, APITimeoutError, RateLimitError
client = AsyncOpenAI(
api_key=os.environ["XYU_API_KEY"],
base_url=os.getenv("XYU_BASE_URL", "https://xyuapi.top/v1"), # 备用 https://xyuai.cc/v1
timeout=60.0,
max_retries=0, # 重试自己写,才能统计真实失败次数
)
MAX_CONCURRENCY = 8
stats = {"ok": 0, "fail": 0}
errors = []
async def ask(sem: asyncio.Semaphore, prompt: str):
"""单条请求:拿到信号量才发,任何异常都在这里消化掉。"""
async with sem:
try:
resp = await client.chat.completions.create(
model="deepseek-v3.2-thinking",
messages=[{"role": "user", "content": prompt}],
max_tokens=512,
)
content = resp.choices[0].message.content
if not content or not content.strip(): # 空回答按失败算,别让 None 流到下游
raise ValueError("empty content")
stats["ok"] += 1
return content.strip()
except RateLimitError as e:
stats["fail"] += 1
errors.append((prompt, "429", str(e)[:200]))
except (APITimeoutError, asyncio.TimeoutError):
stats["fail"] += 1
errors.append((prompt, "timeout", ""))
except Exception as e:
stats["fail"] += 1
errors.append((prompt, type(e).__name__, str(e)[:200]))
return None
async def main(prompts):
if not prompts: # 边界:空输入直接返回,别让 gather 拿到空参数
print("没有待处理任务")
return []
sem = asyncio.Semaphore(MAX_CONCURRENCY)
results = await asyncio.gather(*(ask(sem, p) for p in prompts))
ok_items = [r for r in results if r]
print(f"总数 {len(prompts)} 成功 {stats['ok']} 失败 {stats['fail']}")
if errors:
print("失败样本前三条:", errors[:3])
return ok_items
if __name__ == "__main__":
data = [f"用一句话解释第 {i} 个并发相关的概念" for i in range(1, 21)]
asyncio.run(main(data))
把这段存成 sem_batch.py,设好环境变量 XYU_API_KEY 再 python sem_batch.py 就能跑。接口走 OpenAI 兼容协议,/v1/chat/completions 与 /v1/models 都是标准路径,base_url 指向 https://xyuapi.top/v1 就行;备用地址是 https://xyuai.cc/v1,需要 Gemini 原生协议时另走 /v1beta。
gather 默认遇到异常就往上抛,剩下还在飞的请求结果全部丢掉。上面每个任务自己 try/except,失败返回空值,批次级别的 gather 就永远不会炸。
并发 8 表示同一时刻最多 8 条在飞;RPS 8 表示每秒只放行 8 条新请求。如果每条请求 0.1 秒就返回,并发 8 的实际吞吐是每秒 80 条,早就越过 RPS 上限。只限并发不限 RPS,可能依然被每分钟请求数那道墙拦下。反过来,只限 RPS 不限并发,遇到单条跑 40 秒的慢模型,会同时堆出几百条在途请求。两个都得控。
# rate_limit.py —— 令牌桶 + 与信号量组合使用
import asyncio
import time
from collections import deque
class TokenBucket:
"""按秒放行的限流器:桶最多攒 capacity 个令牌,每秒补 refill_per_sec 个。"""
def __init__(self, capacity: int = 20, refill_per_sec: float = 10.0):
self.capacity = max(1, capacity)
self.refill_rate = max(0.1, refill_per_sec)
self.tokens = float(self.capacity)
self.updated_at = time.monotonic()
self.lock = asyncio.Lock()
self.history = deque(maxlen=200) # 最近放行的时刻,用来观察真实 RPS
def _refill(self):
now = time.monotonic()
delta = now - self.updated_at
if delta <= 0:
return
self.tokens = min(self.capacity, self.tokens + delta * self.refill_rate)
self.updated_at = now
async def acquire(self, need: int = 1):
"""拿到令牌才返回;拿不到就精确睡够时间,不忙等烧 CPU。"""
need = max(1, need)
while True:
async with self.lock:
self._refill()
if self.tokens >= need:
self.tokens -= need
self.history.append(time.monotonic())
return
wait = (need - self.tokens) / self.refill_rate
await asyncio.sleep(min(max(wait, 0.01), 2.0))
def observed_rps(self, window: float = 5.0) -> float:
now = time.monotonic()
recent = [t for t in self.history if now - t <= window]
return len(recent) / window
async def guarded_call(bucket: TokenBucket, sem: asyncio.Semaphore, prompt: str):
"""先过每秒放行这道闸,再去抢并发名额。"""
await bucket.acquire()
async with sem:
try:
resp = await client.chat.completions.create(
model="grok-4.1",
messages=[{"role": "user", "content": prompt}],
max_tokens=256,
)
return (resp.choices[0].message.content or "").strip() or None
except Exception as e:
print(f"失败:{type(e).__name__}")
return None
async def demo():
bucket = TokenBucket(capacity=10, refill_per_sec=5) # 允许短时突发 10 条,长期 5 条/秒
sem = asyncio.Semaphore(8)
tasks = [guarded_call(bucket, sem, f"写一条关于序号 {i} 的短评") for i in range(40)]
results = await asyncio.gather(*tasks)
print("成功", sum(1 for r in results if r), "观测 RPS", round(bucket.observed_rps(), 2))
capacity 是桶里最多能攒多少令牌,决定允许多大的短时突发;refill_per_sec 是长期平均速率。上游限制每分钟 60 条,就把 refill_per_sec 设成 1.0,capacity 设成 5 到 10,留一点突发余量。
guarded_call 里 bucket.acquire() 写在 sem 前面。反过来的话,抢到并发名额的协程会卡在等待上,白占名额,实际吞吐比设定值还低。另外要用 await asyncio.sleep();同步的 time.sleep() 会把事件循环按停,在飞请求一起卡住。
把一万条任务塞进 asyncio.Queue,起 8 个 worker 循环取,取空就退出。worker 数量就是并发上限。这个结构比一次性建一万个协程好在三处:进度实时可打印,失败任务单独落盘,中断重跑按 id 跳过做过的部分。
# pool_batch.py —— 队列 + worker 池,失败任务落盘,重跑只补失败的
import asyncio
import json
import os
from pathlib import Path
from openai import AsyncOpenAI
WORKERS = 8
IN_PATH = Path("input.jsonl")
OUT_PATH = Path("output.jsonl")
FAIL_PATH = Path("failed.jsonl")
client = AsyncOpenAI(
api_key=os.environ["XYU_API_KEY"],
base_url=os.getenv("XYU_BASE_URL", "https://xyuapi.top/v1"),
timeout=120.0,
max_retries=0,
)
counter = {"total": 0, "ok": 0, "fail": 0}
def load_done_ids(path: Path) -> set:
"""读已完成记录的 id,坏行跳过,不因为一行脏数据整批重跑。"""
if not path.exists():
return set()
ids = set()
with path.open(encoding="utf-8") as f:
for line in f:
line = line.strip()
if not line:
continue
try:
ids.add(json.loads(line)["id"])
except (json.JSONDecodeError, KeyError, TypeError):
continue
return ids
async def do_one(item: dict) -> str:
resp = await client.chat.completions.create(
model="deepseek-v3.2-thinking",
messages=[{"role": "user", "content": item["text"]}],
max_tokens=256,
)
text = (resp.choices[0].message.content or "").strip()
if not text:
raise ValueError("empty content")
return text
async def worker(name: int, queue: asyncio.Queue, sem: asyncio.Semaphore,
out_f, fail_f):
while True:
try:
item = queue.get_nowait() # 队列空就退出,不等死
except asyncio.QueueEmpty:
return
try:
async with sem:
text = await asyncio.wait_for(do_one(item), timeout=120)
out_f.write(json.dumps({"id": item["id"], "text": text},
ensure_ascii=False) + "\n")
out_f.flush()
counter["ok"] += 1
except asyncio.TimeoutError:
counter["fail"] += 1
fail_f.write(json.dumps({"id": item["id"], "text": item["text"],
"err": "timeout"}, ensure_ascii=False) + "\n")
fail_f.flush()
except Exception as e:
counter["fail"] += 1
fail_f.write(json.dumps({"id": item["id"], "text": item["text"],
"err": f"{type(e).__name__}: {e}"[:200]},
ensure_ascii=False) + "\n")
fail_f.flush()
finally:
queue.task_done()
done = counter["ok"] + counter["fail"]
if done % 500 == 0:
print(f"进度 {done}/{counter['total']} 成功 {counter['ok']} 失败 {counter['fail']}")
async def main():
if not IN_PATH.exists():
print(f"找不到 {IN_PATH},先准备输入文件")
return
done_ids = load_done_ids(OUT_PATH)
items = []
with IN_PATH.open(encoding="utf-8") as f:
for line in f:
line = line.strip()
if not line:
continue
try:
obj = json.loads(line)
except json.JSONDecodeError:
continue
if "id" not in obj or "text" not in obj:
continue
if obj["id"] in done_ids:
continue
items.append(obj)
if not items:
print("没有待补的任务,全部已完成")
return
counter["total"] = len(items)
queue = asyncio.Queue()
for it in items:
queue.put_nowait(it)
sem = asyncio.Semaphore(WORKERS)
with OUT_PATH.open("a", encoding="utf-8") as out_f, \
FAIL_PATH.open("a", encoding="utf-8") as fail_f:
await asyncio.gather(*(worker(i, queue, sem, out_f, fail_f)
for i in range(WORKERS)))
print(f"结束:成功 {counter['ok']} 失败 {counter['fail']},"
f"失败明细见 {FAIL_PATH}")
if __name__ == "__main__":
asyncio.run(main())
输入文件每行一个 JSON,字段是 id 和 text。worker 用 get_nowait() 加 QueueEmpty 退出,比等 queue.join() 省心,任务数少于 worker 数时也不卡住。
失败明细写进 failed.jsonl,带 id、原文和错误类型。补跑时把它改名成 input.jsonl 再跑一遍,成功记录按 id 去重,不会重复计费。output.jsonl 用追加模式打开,进程被杀时已经 flush 的行依然有效,断电不用从零开始。
ThreadPoolExecutor 和结果顺序老项目里全是同步 SDK,改成异步要动一大片代码,这时用 ThreadPoolExecutor(max_workers=8) 最省事:八条线程轮流发请求,同样把在途请求压在 8 条以内。每条请求还可能带自己的重试,实际瞬时请求数会更高一点,起步值别开太大。
as_completed 和 map 的顺序差别pool.map(ask, prompts) 返回的结果顺序和输入严格一致,内部照样并发执行,适合结果要按原顺序写回文件的场景。as_completed(futures) 是谁先返回谁先出,顺序随机,适合只要结果集合的场景。两条路都能保序:as_completed 配合 {future: index} 的记录,收完按 index 排回来就行。
# sync_batch.py —— 同步 SDK 走线程池,两种收结果的方式都在
import os
import random
import time
from concurrent.futures import ThreadPoolExecutor, as_completed, TimeoutError
from openai import APIError, OpenAI
client = OpenAI(
api_key=os.environ["XYU_API_KEY"],
base_url=os.getenv("XYU_BASE_URL", "https://xyuapi.top/v1"),
timeout=60.0,
max_retries=0,
)
def ask(prompt: str) -> str:
"""单条调用,最多三次,指数退避加抖动,失败返回标记字符串。"""
if not prompt or not prompt.strip():
return "__FAIL__ empty prompt"
for attempt in range(3):
try:
resp = client.chat.completions.create(
model="grok-4.1",
messages=[{"role": "user", "content": prompt}],
max_tokens=256,
)
text = (resp.choices[0].message.content or "").strip()
return text if text else "__FAIL__ empty content"
except APIError as e:
if attempt == 2:
return f"__FAIL__ {type(e).__name__}"
time.sleep(2 ** attempt + random.random())
return "__FAIL__ unknown"
def run_ordered(prompts):
"""map:输出顺序与输入一致,结果可以直接按行写回文件。"""
with ThreadPoolExecutor(max_workers=8) as pool:
return list(pool.map(ask, prompts))
def run_fast(prompts, batch_timeout=600):
"""as_completed:谁先完谁先出,缺的用失败标记补齐,保持长度一致。"""
out = {}
with ThreadPoolExecutor(max_workers=8) as pool:
futures = {pool.submit(ask, p): i for i, p in enumerate(prompts)}
try:
for fut in as_completed(futures, timeout=batch_timeout):
idx = futures[fut]
try:
out[idx] = fut.result()
except Exception as e:
out[idx] = f"__FAIL__ {type(e).__name__}"
except TimeoutError:
print(f"超时:还有 {len(prompts) - len(out)} 条没返回,按失败记录")
return [out.get(i, "__FAIL__ timeout") for i in range(len(prompts))]
if __name__ == "__main__":
prompts = [f"给编号 {i} 的条目写一句说明" for i in range(1, 51)]
ordered = run_ordered(prompts)
fast = run_fast(prompts)
print("顺序一致:", ordered[0][:20])
print("失败条数:", sum(1 for r in fast if r.startswith("__FAIL__")))
Retry-After:服务端告诉你还要等几秒。有它就直接按它等,不要自己猜。x-ratelimit-remaining-requests:当前窗口还剩多少请求额度,快见底时主动降速。x-ratelimit-reset-requests:额度重置的时间,可能是一个秒数,也可能是 12s 这种带单位的串。限流响应头属于可选信息,不同上游、不同版本可能给也可能不给,429 响应里干脆一个都没有的情况很常见。代码里必须写兜底路径:先试 retry-after,没有就试重置时间,再没有就用默认的指数退避加抖动。只按响应头写逻辑的脚本,换一个入口就会原地空转。
退避基准用 1 到 2 秒,每次翻倍,最高压到 30 秒封顶,重试次数不超过 3 次。抖动是给等待时间乘一个 1 到 2 之间的随机系数,用来打散同时回头的请求。没有抖动,几百条请求会在同一毫秒一起重发,再撞一次上游。
固定并发值最大的问题是不知道上游今天的状态:上午 8 并发一把过,下午上游负载高了,8 也开始 429。AIMD 的思路很简单:连续 50 次请求成功,并发加一档;撞到一次 429,立刻砍一半,并把踩线的那个值记下来。涨得慢、退得快,既吃满上游容量,也不容易持续撞墙。
# adaptive.py —— AIMD 自适应并发,撞 429 减半,连续成功慢慢加回来
import asyncio
import random
from contextlib import asynccontextmanager
from openai import APITimeoutError, RateLimitError
def retry_after_seconds(exc, default: float = 1.5) -> float:
"""优先读 Retry-After,其次读重置时间,都读不到就用默认值加抖动。"""
resp = getattr(exc, "response", None)
headers = getattr(resp, "headers", {}) if resp is not None else {}
for key in ("retry-after", "x-ratelimit-reset-requests"):
raw = headers.get(key)
if not raw:
continue
try:
return max(0.5, float(str(raw).strip().rstrip("s")))
except ValueError:
continue
return default * (1 + random.random())
def rate_limit_left(exc):
"""剩余额度,取不到返回 None,调用方按未知处理。"""
resp = getattr(exc, "response", None)
headers = getattr(resp, "headers", {}) if resp is not None else {}
for key in ("x-ratelimit-remaining-requests", "x-ratelimit-remaining"):
if key in headers:
try:
return int(headers[key])
except (TypeError, ValueError):
return None
return None
class AdaptiveLimiter:
"""用条件变量放行:in-flight 达到 limit 就在门口等,降档时新请求自动排队。"""
def __init__(self, start: int = 4, ceil: int = 32, floor: int = 1, step_ok: int = 50):
self.limit = max(floor, min(ceil, start))
self.ceil = ceil
self.floor = floor
self.step_ok = step_ok
self.ok_streak = 0
self.known_bad = None # 踩过线的并发值,重启时可以直接从它减半起步
self.inflight = 0
self.cond = asyncio.Condition()
@asynccontextmanager
async def slot(self):
async with self.cond:
await self.cond.wait_for(lambda: self.inflight < self.limit)
self.inflight += 1
try:
yield
finally:
async with self.cond:
self.inflight -= 1
self.cond.notify_all()
async def _scale(self, up: bool):
async with self.cond:
if up:
new = min(self.ceil, self.limit + 1)
else:
self.known_bad = self.limit
new = max(self.floor, self.limit // 2)
if new != self.limit:
print(f"并发调整 {self.limit} -> {new}")
self.limit = new
self.cond.notify_all()
async def call(self, prompt: str, max_retries: int = 3):
for attempt in range(max_retries):
try:
async with self.slot():
resp = await client.chat.completions.create(
model="deepseek-v3.2-thinking",
messages=[{"role": "user", "content": prompt}],
max_tokens=256,
)
async with self.cond:
self.ok_streak += 1
bump = self.ok_streak >= self.step_ok
if bump:
self.ok_streak = 0
if bump:
await self._scale(up=True)
return (resp.choices[0].message.content or "").strip()
except RateLimitError as e:
wait = retry_after_seconds(e, default=1.5 * (2 ** attempt))
await self._scale(up=False)
print(f"撞 429,等 {wait:.1f} 秒后第 {attempt + 1} 次重试")
await asyncio.sleep(wait)
except APITimeoutError:
await asyncio.sleep(1.0 + random.random())
return None
降档不会打断已经在飞的请求:inflight 大于新 limit 时,新请求在门口多等一会儿,等旧请求还完名额再放行。只让入口变窄,不加锁强杀。
起步值取 4,别嫌少。每成功 50 次请求把并发加一,撞到一次 429 立刻退回一半,并把退回前的数值当作上限记住,下次从「记住的值的一半」起步,不要再从 4 开始摸索。
下面这组数字是本地一万条任务的示例,单条平均 1.8 秒,只用来说明趋势,不代表任何上游的硬指标。
| 并发数 | 请求成功率 | 总耗时 | 实际情况 |
|---|---|---|---|
| 1 | 100% | 约 5.0 小时 | 稳,时间全花在等 |
| 4 | 100% | 约 1.3 小时 | 起步档,几乎无 429 |
| 8 | 99% | 约 40 分钟 | 偶发 429,退避能吞掉 |
| 16 | 88% | 约 55 分钟 | 退避开始拉长总时间 |
| 32 | 61% | 约 1.6 小时 | 大量打回,越努力越慢 |
| 64 | 34% | 约 3.0 小时 | 基本在互相踩,纯浪费 |
这张表最值得看的是 16 那一行:成功率只掉 12 个百分点,总耗时却比 8 并发更久。失败的请求要等退避再重发,等待时间算进总时长。
限流的职责是不让在途请求数超过上游能吃下的量,作用在请求发出之前。重试的职责是某条请求真失败了把它捞回来,作用在失败之后。两者缺一个都不行:只有限流没有重试,一次网络抖动就让你丢数据;只有重试没有限流,重试自己会变成新的流量洪峰。限流决定发多快,重试决定失败了怎么办;重试的等待时间要按响应头或退避算出来,不能原地立刻重发。
| 响应或错误 | 含义 | 该做什么 |
|---|---|---|
| 429 Too Many Requests | 并发或频率超额度 | 立刻降并发,按 Retry-After 等待 |
| 500 / 502 / 503 | 上游内部抖动 | 退避 2 秒起重试,最多三次 |
| 请求超时 | 网络慢或长文本输出 | 把超时提到 120 秒,再考虑重试 |
| 连接被重置 | 长连接被中间链路掐断 | 重建连接重试同一条 |
| 401 Unauthorized | 令牌无效或已过期 | 不要重试,检查令牌与 base_url |
| 400 Bad Request | 请求参数不合法 | 不要重试,先修请求体 |
| 404 model not found | 模型名写错或已下线 | 不要重试,先查模型列表 |
按次计费是主流方式:一次请求一个固定价,输入写多长都不影响单价,配合按量计费和无限卡套餐一起用。deepseek-v3.2-thinking 是 0.049 元一次,5000 次顺利跑完就是 5000 乘 0.049,等于 245 元。充值门槛不高,最低 7 元,支付宝或微信直接付,不需要海外信用卡。
不控并发的那一版账是这样的:大约四成请求被 429 打回,也就是 2000 条要重发,再算 2000 乘 0.049,多出 98 元,总支出变成 343 元。退避等待还让整体耗时接近翻倍,40 分钟的活儿拖到一个半小时。把并发压住、按服务端要求的时间退避,2000 次重试直接省掉,245 元跑完,时间回到 40 分钟那一档。
| 模型 | 单价(元/次) |
|---|---|
| deepseek-v3.2-thinking | 0.049 |
| deepseek-v4-flash-thinking | 0.05 |
| grok-4.1 | 0.05 |
| claude-sonnet-4-5-thinking | 0.09 |
| gpt-5.5 | 0.2 |
| claude-opus-4-6-thinking | 0.25 |
写脚本时把这三件事分开:并发用信号量摁住,频率用令牌桶卡住,失败交给带退避的重试捞。跑之前先用 20 条数据试一轮,看一眼成功率和实际耗时,再决定要不要把并发往上抬。多花五分钟做这件事,能省下来的是一整晚的重跑和一笔白花的钱。