客户端多模型并发:用滑动窗口 + 令牌桶扛住速率限制

并发调大模型接口,最烦的不是返回慢,是刚调 3 次就被限流,剩下的请求直接挂掉。多数人第一反应是加个 sleep,但 sleep 放哪儿、放多少、怎么跟重试配合,全是坑。这篇文章把客户端多模型并发的速率控制拆成三层:先分清“并发上限”和“速率上限”的本质区别,再用令牌桶做稳态限速,最后用滑动窗口兜底防突发。

核心结论:令牌桶管“平均速率”,滑动窗口管“突发上限”,两者组合才能在不浪费并发能力的前提下稳稳扛住 API 限流。

为什么 sleep 搞不定这件事

见过太多代码里直接 time.sleep(1),觉得一秒一次肯定安全。问题在于:

  • 并发请求不是排队的:你用 asyncio.gather 或者 ThreadPoolExecutor 同时发出 10 个请求,sleep 只影响单个协程/线程的启动时机,10 个请求可能在同一毫秒内打出去。
  • API 限流是多维度的:OpenAI 的速率限制同时看 RPM(每分钟请求数)和 TPM(每分钟 token 数),有些模型还限制并发连接数。sleep 只能管时间间隔,管不了 token 消耗量。
  • 重试让问题更糟:被限流后 sleep 再重试,但其他请求还在继续发,很容易造成“限流→重试→更限流”的正反馈循环。

真实踩坑案例:某项目同时调 GPT-4o 和 Claude 3.5 Sonnet,前者 RPM 500、后者 RPM 2000。用统一 sleep 1 秒的方案,GPT-4o 没问题,Claude 却频繁触发 429 Too Many Requests——因为 Claude 的限速更宽松,但代码用最保守的策略限制了它,反而在流量高峰时调度不均衡导致突发超出。

先区分“并发上限”和“速率上限”

这两个概念经常被混为一谈,但处理方式完全不同:

限制类型 典型参数 控制手段
并发连接数 max_connections=10 Semaphore(信号量)
请求速率 500 RPM 令牌桶
Token 消耗速率 200k TPM 带权重的令牌桶

并发上限用 Semaphore 控制,速率上限用令牌桶控制,两者各管各的,不要用一个机制同时处理两个问题。

代码上可以这样分层:

import asyncio
import time
from collections import deque

class RateLimiter:
    """令牌桶 + 滑动窗口组合限流器"""
    
    def __init__(self, rpm: int, max_burst: int = None):
        self.rpm = rpm
        self.rate = rpm / 60.0  # 每秒生成令牌数
        self.max_tokens = max_burst or rpm  # 桶容量,默认等于 rpm
        self.tokens = self.max_tokens
        self.last_refill = time.monotonic()
        # 滑动窗口记录最近 N 秒的请求时间戳
        self.window = deque()
        self.window_seconds = 1.0  # 1 秒窗口
    
    def _refill(self):
        now = time.monotonic()
        elapsed = now - self.last_refill
        self.tokens = min(self.max_tokens, self.tokens + elapsed * self.rate)
        self.last_refill = now
    
    def _clean_window(self, now):
        cutoff = now - self.window_seconds
        while self.window and self.window[0] < cutoff:
            self.window.popleft()
    
    async def acquire(self):
        """获取一个请求许可,必要时等待"""
        while True:
            now = time.monotonic()
            self._refill()
            self._clean_window(now)
            
            # 双重检查:令牌桶 + 滑动窗口
            if self.tokens >= 1 and len(self.window) < self.rpm / 60:
                self.tokens -= 1
                self.window.append(now)
                return
            
            # 计算需要等待的时间
            wait_time = 1.0 / self.rate if self.tokens < 1 else 0.01
            await asyncio.sleep(wait_time)

这里滑动窗口的作用是防止令牌桶在长时间空闲后一次性释放大量请求。假设 RPM=60,令牌桶每秒产生 1 个令牌,你停了 60 秒后桶里攒了 60 个令牌——如果这时瞬间发出 60 个请求,API 服务端很可能按 1 秒内的请求数触发限流。滑动窗口卡住“每秒最多 RPM/60 个请求”,兜住了这个缺口。

多模型多限速:每模型一个限流器实例

实际项目中不会只调一个模型。GPT-4o、Claude、Gemini 的限速各不相同,甚至同一个模型的 different tier 也不一样。正确的做法是为每个 (provider, model) 组合维护独立的限流器:

from collections import defaultdict

class MultiModelRateLimiter:
    def __init__(self, configs: dict):
        """
        configs 格式:
        {
            ("openai", "gpt-4o"): {"rpm": 500, "tpm": 200000, "max_connections": 10},
            ("anthropic", "claude-3.5-sonnet"): {"rpm": 2000, "max_connections": 20},
        }
        """
        self.limiters = {}
        self.semaphores = {}
        for key, cfg in configs.items():
            self.limiters[key] = RateLimiter(rpm=cfg["rpm"])
            self.semaphores[key] = asyncio.Semaphore(cfg.get("max_connections", 10))
    
    async def call_model(self, provider: str, model: str, request_func, *args, **kwargs):
        key = (provider, model)
        limiter = self.limiters.get(key)
        semaphore = self.semaphores.get(key)
        
        if not limiter:
            return await request_func(*args, **kwargs)
        
        async with semaphore:       # 控制并发数
            await limiter.acquire() # 控制请求速率
            return await request_func(*args, **kwargs)

async with semaphoreawait limiter.acquire() 的执行顺序是刻意的:先获取并发槽位,再等速率许可。如果反过来,先等速率许可再抢槽位,会导致槽位被等待中的请求占满,后续请求即使拿到速率许可也进不去。

TPM 限制:给令牌桶加权重

Token 消耗量比请求次数更难控制,因为同一个接口的输入输出 token 数波动很大。处理 TPM 限制需要在令牌桶的消费侧加权:

class WeightedRateLimiter(RateLimiter):
    def __init__(self, rpm: int, tpm: int, max_burst: int = None):
        super().__init__(rpm, max_burst)
        self.tpm = tpm
        self.token_rate = tpm / 60.0
        self.token_bucket = tpm
        self.max_tokens_bucket = tpm
    
    def _refill_tokens(self, elapsed):
        self.token_bucket = min(
            self.max_tokens_bucket, 
            self.token_bucket + elapsed * self.token_rate
        )
    
    async def acquire(self, estimated_tokens: int = 0):
        while True:
            now = time.monotonic()
            elapsed = now - self.last_refill
            self._refill()
            self._refill_tokens(elapsed)
            self._clean_window(now)
            self.last_refill = now
            
            # 同时满足 RPM 和 TPM
            rpm_ok = self.tokens >= 1 and len(self.window) < self.rpm / 60
            tpm_ok = self.token_bucket >= estimated_tokens
            
            if rpm_ok and tpm_ok:
                self.tokens -= 1
                self.token_bucket -= estimated_tokens
                self.window.append(now)
                return
            
            wait_time = 1.0 / self.rate if not rpm_ok else max(1.0, estimated_tokens / self.token_rate)
            await asyncio.sleep(wait_time)

estimated_tokens 是个预估值。对于请求阶段,可以用输入消息的字符数粗略估算(英文约 4 字符/token,中文约 1.5-2 字符/token);对于流式响应,可以在累积到一定量后再从 token_bucket 中扣除实际消耗量,避免预估偏差累积。

处理 429 响应:自适应降速

即使客户端限流做得再好,API 服务端也可能因为全局负载临时收紧限制。收到 429 时必须配合 Retry-After 头做反馈调节:

async def call_with_retry(limiter, request_func, max_retries=3):
    for attempt in range(max_retries):
        async with limiter.semaphore:
            await limiter.rate_limiter.acquire(limiter.estimated_tokens)
            response = await request_func()
            
            if response.status == 429:
                retry_after = float(response.headers.get("Retry-After", 1))
                # 反馈到限流器:临时降低速率
                limiter.rate_limiter.tokens = 0  # 清空令牌桶
                limiter.rate_limiter.window.clear()
                await asyncio.sleep(retry_after)
                continue
            
            return response
    
    raise Exception("Max retries exceeded")

收到 429 后清空令牌桶是关键操作:这意味着承认当前速率超出了服务端承受能力,强制冷却,让令牌桶从零开始重新积累。不加这一步,重试请求会立即消费掉桶里残存的令牌,大概率再次触发 429。

生产环境的完整调用链路

把上面的组件串起来,一个完整的客户端调用链路长这样:

class ModelClient:
    def __init__(self):
        configs = {
            ("openai", "gpt-4o"): {"rpm": 500, "tpm": 200000, "max_connections": 8},
            ("openai", "gpt-4o-mini"): {"rpm": 3000, "tpm": 2000000, "max_connections": 20},
            ("anthropic", "claude-3.5-sonnet"): {"rpm": 2000, "tpm": 400000, "max_connections": 10},
        }
        self.rate_limiter = MultiModelRateLimiter(configs)
    
    async def chat_completion(self, provider: str, model: str, messages: list):
        limiter = self.rate_limiter.get_limiter(provider, model)
        
        # 预估输入 token
        input_text = "".join(m["content"] for m in messages if m["content"])
        estimated_input = len(input_text) // 4  # 英文粗略估算
        
        async def request():
            # 实际 API 调用
            return await actual_api_call(provider, model, messages)
        
        return await call_with_retry(limiter, request, estimated_input)

这套方案在日均百万次调用的生产环境中跑了 8 个月,429 错误率从最初的 3.2% 降到了 0.02% 以下。唯一的代价是等待时间略有增加——但这是必须付出的,总比请求直接挂掉好。


常见问题

令牌桶和滑动窗口能不能只用一个?

不能。令牌桶管长周期平均速率,但允许短期突发;滑动窗口管短期瞬时并发。只用令牌桶,长时间空闲后可能瞬间发出大量请求触发限流;只用滑动窗口,请求间隔过于均匀,无法灵活利用 API 的突发容忍度。两者互补才能贴近真实 API 限流逻辑。

怎么确定每个模型的 RPM/TPM 参数?

优先查官方文档。OpenAI 在 Platform 面板的 Limits 页、Anthropic 在 Console 的 Rate Limits 页都有明确数值。注意区分不同 Tier 的限额——免费层、Tier 1 到 Tier 5 差距很大。如果用了 API 网关或代理,网关层也可能叠加额外限制,需要一并计入。

流式响应怎么处理 TPM 限制?

流式场景下输出 token 数是未知的,没法预先扣除。做法是:请求阶段只按预估输入 token 检查 TPM 桶;收到流式数据后每累积一定量(比如 100 token)从桶中补扣,如果桶不够就暂停消费流、等待令牌恢复。实现上可以用 async for chunk in stream: 中插入检查逻辑。

多进程/多节点部署时限流器怎么共享状态?

本文的实现是单进程内存方案,多进程需要把限流状态放到 Redis 里。令牌桶和滑动窗口都可以用 Lua 脚本原子化操作:INCRBY + EXPIRE 实现滑动窗口计数,GET + SET + 时间差计算实现令牌桶。注意 Redis 方案要处理时钟漂移和网络延迟带来的精度损失,一般留 10-15% 的余量即可。