用令牌桶算法在客户端主动限流,避免 429 风暴;60 行代码同时处理请求数和 token 数两个维度,并支持多租户公平性。
速率限制通常是两个数字——每分钟请求数和每分钟 token 数,两者有其一超限就算突破。一个单独的 token bucket 可以同时处理这两者,大约六十行代码,就能把 429 风暴转化为一个只比限制慢几毫秒的队列。
重试 429 是有效的。但不触发它更好,原因有四个具体理由。
一个被拒绝的请求仍然消耗时间。往返、退避、再往返。客户端等待 200ms 就能替代一次失败加两秒退避。
重试会把一次超限转化为多次。每一次重试的请求都会在恰好超限的时刻增加负载,这就是为什么速率限制事件会升级而非平息。
有些限制是惩罚性的。一个通过延长冷却时间来应对持续超限的提供商,会让一个自我造成的问题持续远超突发时段。
限速器是一个放置公平性的地方。每个租户一个 bucket 可以阻止单个客户的批处理作业耗尽整个 key,这是任何重试策略都做不到的。
整个算法不过是一行算术。桶最多容纳 capacity 个单位,并以 rate 个单位每秒的速度补充。要花费 n 个单位,你需要等待直到至少存在 n 个,然后减去它们。
available(now) = min(capacity, available(last) + rate * (now - last))
要花费 n:
如果 available >= n:立即花费
否则:等待 (n - available) / rate 秒,然后花费
对于每分钟 3000 次请求的限制:
rate = 3000 / 60 = 50 每秒请求数
capacity = 50 (一秒的突发)或者 3000(一分钟的突发)
capacity 的选择是设计中唯一的判断,而且这是一个真实的权衡。capacity 等于一秒的 rate 意味着空闲客户端完全无法突发——平滑、安全,但从队列中排空的速度较慢。capacity 等于整个窗口允许一个空闲了一分钟的客户端立即发出整个分钟的配额,这恰恰是在提供商的限速器(如果他们是在更短的间隔内测量的话)上触发的那种突发。在一秒和十秒的 rate 之间是一个可辩护的默认值;从一秒开始,如果吞吐量受损再提高。
注意算法中没有定时器、没有后台线程、也没有队列。桶的容量是通过流逝时间来计算的,这就是为什么实现如此简短,以及为什么它不会漂移。
# limiter.py
import asyncio
import threading
import time
class TokenBucket:
"""A refilling bucket. Thread-safe; see AsyncTokenBucket for asyncio."""
def __init__(self, rate_per_second: float, capacity: float | None = None):
if rate_per_second <= 0:
raise ValueError("rate must be positive")
self.rate = float(rate_per_second)
self.capacity = float(capacity if capacity is not None else rate_per_second)
self._available = self.capacity
self._updated = time.monotonic()
self._lock = threading.Lock()
def _refill(self) -> None:
now = time.monotonic()
elapsed = now - self._updated
if elapsed > 0:
self._available = min(self.capacity, self._available + elapsed * self.rate)
self._updated = now
def _wait_time(self, amount: float) -> float:
"""Reserve the amount, returning how long the caller must sleep first."""
if amount > self.capacity:
raise ValueError(
f"cannot spend {amount}: bucket capacity is {self.capacity}"
)
with self._lock:
self._refill()
self._available -= amount # may go negative: that IS the queue
deficit = -self._available
return max(0.0, deficit / self.rate)
def acquire(self, amount: float = 1.0) -> None:
delay = self._wait_time(amount)
if delay > 0:
time.sleep(delay)
def refund(self, amount: float) -> None:
"""Give back units reserved but not used."""
with self._lock:
self._refill()
self._available = min(self.capacity, self._available + amount)
关键的设计选择是 _available 允许变成负数。在锁内减去并在锁外睡眠意味着每个调用者从一个隐式队列中的位置计算出自己的等待时间,因此二十个同时请求的线程会得到二十个交错唤醒,而不是在共享睡眠后二十个同时唤醒。这个雷鸣群效应(thundering-herd)bug 存在于网上大多数简陋的限速器中,而且它产生的正是这个限速器本应阻止的突发。
使用 time.monotonic() 而不是 time.time(),原因与日志方案中相同:时钟调整不能让你获得一分钟的免费容量或让桶停滞一小时。
asyncio 版本是同样的算术,只是睡眠不同:
class AsyncTokenBucket(TokenBucket):
def __init__(self, rate_per_second: float, capacity: float | None = None):
super().__init__(rate_per_second, capacity)
self._alock = asyncio.Lock()
async def acquire(self, amount: float = 1.0) -> None: # type: ignore[override]
async with self._alock:
self._refill()
self._available -= amount
delay = max(0.0, -self._available / self.rate)
if delay > 0:
await asyncio.sleep(delay)
不要在线程和 asyncio 任务之间共享一个 TokenBucket。同步版本使用 time.sleep 睡眠,这会阻塞整个事件循环——这正是用 asyncio 同时发出四十个请求时描述的那种失败。二选其一,每个进程一个。
每分钟请求数和每分钟 token 数是两个桶。一次调用必须同时满足两者,而等待时间是两者中较长的——如果你简单地依次从每个桶获取,这自然就出来了。
# governed.py
class ProviderLimits:
def __init__(self, rpm: int, tpm: int, burst_seconds: float = 1.0):
self.requests = TokenBucket(rpm / 60.0, capacity=rpm / 60.0 * burst_seconds)
self.tokens = TokenBucket(tpm / 60.0, capacity=tpm / 60.0 * burst_seconds)
def acquire(self, estimated_tokens: int) -> None:
self.requests.acquire(1)
self.tokens.acquire(estimated_tokens)
def reconcile(self, estimated_tokens: int, actual_tokens: int) -> None:
difference = estimated_tokens - actual_tokens
if difference > 0:
self.tokens.refund(difference) # we over-reserved
elif difference < 0:
self.tokens.acquire(-difference) # we under-reserved: pay it back
limits = ProviderLimits(rpm=3_000, tpm=1_000_000)
def governed_call(client, payload: dict, estimated_tokens: int) -> dict:
limits.acquire(estimated_tokens)
response = client.post("/chat/completions", json=payload)
response.raise_for_status()
body = response.json()
usage = body.get("usage", {})
actual = usage.get("total_tokens")
if actual is not None:
limits.reconcile(estimated_tokens, actual)
return body
先获取请求槽再获取 token 槽是刻意的:请求槽更便宜,而且即使 token 估计很差,它也能保持请求计数精确。
你必须在调用之前预留容量,但真实 token 数只能在调用之后才知道。输入 token 可以提前知道;输出 token 不能,而且 max_tokens 和模型实际生成的内容之间的差距可能达到二十倍。
预留 max_tokens 是安全但浪费的——一个为平均只有 200 token 的回答预留 4000 输出 token 的限速器,会把你限制在实际配额的二十分之一。预留一个估计值并在之后 reconciliation,如上所示,保持长期平均正确,同时允许短期超调。这几乎在所有情况下都是正确的权衡,因为提供商的限制是在一个窗口内强制执行的,而不是瞬时的。
# 调用前的一个可用估计
def estimate_tokens(payload: dict) -> int:
chars = sum(len(m.get("content") or "") for m in payload["messages"])
prompt = chars // 4 # ~4 chars/token for English prose
completion = min(payload.get("max_tokens", 512), 512)
return prompt + completion
每 token 四个字符的规则是英文的一个粗略近似值,对于代码、其他语言以及任何有大量标点符号的内容都明显错误。每词的 token 数和 tokeniser 语言税可以量化误差。要准确就运行实际的 tokeniser;对于之后会 reconciliation 的限速器,这个近似值就够了,运行一天真实流量后你自己的日志会给你一个比任何经验法则都更好的乘数。
从仪表板中的数字配置你的限速器,意味着在账户升级当天它就错了。大多数提供商在每次响应中报告实时限制,所以限速器可以学习它。
响应头名称没有标准化。广泛使用的约定是以 x-ratelimit- 开头的家族,带有 -limit、-remaining 和 -reset 后缀,有时会分成请求和 token 变体;一些网关使用 IETF 草案中的 RateLimit 家族,有些则什么都不发送。防御性地读取它们,在根据一个名称写代码之前先打印你的端点实际返回的内容:
# 一次,从 REPL 中,看看你实际得到了什么
response = client.post("/chat/completions", json=payload)
for name, value in response.headers.items():
if "ratelimit" in name.lower() or name.lower() == "retry-after":
print(f"{name}: {value}")
def observe_headers(headers, limits: ProviderLimits) -> None:
"""Adjust the local buckets from whatever the provider reported."""
remaining = headers.get("x-ratelimit-remaining-requests")
if remaining is not None:
try:
free = float(remaining)
except ValueError:
return
# never grant more than the provider says is left
with limits.requests._lock:
limits.requests._refill()
limits.requests._available = min(limits.requests._available, free)
只向下调整是安全的方向,这是这里整个设计规则。一个报告比你桶认为的剩余空间更多的响应头可能是过期的、可能是在计算不同的窗口、或者可能属于共享网关上的不同 key——根据它提高你的本地配额会把一个正常工作的限速器变成一个间歇性的限速器。降低永远是安全的。
唯一值得无条件执行的头是 429 上的 Retry-After,这是一个直接指令而不是估计;tenacity 中重试模型调用时的 wait 函数已经遵循它。即使你不对剩余配额值做任何操作,也要将它们与你的调用一起记录——一个在一周内趋向于零的剩余空间是你在 429 开始之前收到的最早警告。
上面的 bucket 是每个进程的。四个 rpm=3000 的 uvicorn worker 每分钟会发送 12,000 个请求,而限速器会报告一切正常。
分配预算。零基础设施的答案:给 N 个进程每个 rpm/N。正确,但当负载不均匀时是浪费的——空闲 worker 的份额是不可用的。
Redis 中一个共享 bucket。用 Lua 脚本实现同样的算术,这样 refill-and-subtract 是原子的,并将 available 和 updated 存储在一个 hash 的字段上。往返延迟大约一毫秒,相比模型调用什么都不是。
把它放在所有东西前面。所有进程都调用的网关或代理是唯一一种添加第五个服务不需要重新分配任何人预算的设计,而且这也是每个租户限速自然存在的地方。
无论你选择哪种,都在 tenacity 重试模型调用的重试策略下面保持它。限速器减少 429;它不能消除它们,因为你看到的限制总是略微过时,提供商可能在计算一个你看不到的窗口。Rate limits explained 涵盖了响应头告诉你的关于那个窗口的信息。
Forty Requests at Once With asyncio
Retrying Model Calls With tenacity
Logging Every Model Call