跳转至

如何根据 API 的 RPM TPM 限额合理设置 Semaphore 并发数,并结合 tenacity 构建完整并发控制方案?

这个问题很经典,调用外部 API 时几乎都会碰到。我分享一下我的思考框架和落地方式,带着具体的代码思路,既讲原理也讲坑。


🔍 1. 先搞清楚两个概念:RPM 与 TPM

  • RPM (Requests Per Minute):每分钟最多发多少次请求。

  • TPM (Tokens Per Minute):每分钟最多消耗多少个 token(常见于 LLM API,如 OpenAI)。

你的并发方案必须同时满足两个限制,缺一不可。


📐 2. 从限额推导合理的 Semaphore 并发数

Semaphore 控制的是“同一时刻最多有多少个请求在飞”。核心逻辑:一个槽位在 60 秒内能完成多少次请求?

假设:

  • 每个请求平均耗时 D 秒(从发出到收到完整响应)

  • 每个请求平均消耗 T 个 token

那么:

一个槽位每分钟能完成的请求数 = 60 / D
总请求数 ≤ RPM  ⇒  并发数 * (60 / D) ≤ RPM
                  并发数 ≤ RPM * D / 60

总token消耗 ≤ TPM  ⇒  并发数 * (60 / D) * T ≤ TPM
                  并发数 ≤ TPM * D / (60 * T)

👉 最终安全并发数:

semaphore_size = min( RPM * D / 60 ,  TPM * D / (60 * T) ) * 0.8

乘以 0.8 是安全裕度,防止网络抖动、突发流量或者请求耗时分布不均。

举一个具体的例子:

  • RPM=500,TPM=200000,D=2s,T=1000

  • RPM * D / 60 = 500 * 2 / 60 ≈ 16.67

  • TPM * D / (60 * T) = 200000 * 2 / (60 * 1000) ≈ 6.67

  • 取较小值 6.67,再乘 0.8 → 并发数设为 5 比较稳。

如果没有 TPM 限制,就只看 RPM 那条公式。


🔒 3. 用 Semaphore 实现“在飞”并发控制(asyncio 版)

import asyncio

semaphore = asyncio.Semaphore(5)  # 按计算出的值动态设置

async def call_api(payload):
    async with semaphore:
        # 实际 HTTP 请求
        response = await client.post(url, json=payload)
        return response

这里的关键是 async with semaphore 会保证:

  • 超过限制的协程会被挂起等待

  • 请求完成(无论成功还是抛异常)后自动释放槽位,避免死锁


🔁 4. tenacity 负责“被限流后的重试”

Semaphore 只控制并发数量,不能避免 429,因为你的速率可能刚好卡在临界点,或者有瞬时尖峰。这时就轮到 tenacity 出场,专门处理限流退避。

from tenacity import (
    retry, stop_after_attempt, wait_exponential, retry_if_result
)

def is_rate_limited(response):
    return response.status_code == 429

@retry(
    retry=retry_if_result(is_rate_limited),
    wait=wait_exponential(multiplier=1, min=2, max=60),
    stop=stop_after_attempt(5),
    reraise=True
)
async def api_call_with_retry(payload):
    async with semaphore:            # 重试期间继续占用信号量,保证这个请求
        response = await client.post(url, json=payload)
        if response.status_code == 429:
            # 可选:解析 Retry-After 头,动态指导等待时间
            retry_after = response.headers.get("Retry-After")
            # tenacity 的 wait 会自行等待,这里只需要让函数判定为需要重试
        return response

设计要点:

  • 重试时信号量保持不变:一次请求的多次重试共用同一个并发槽位,防止雪崩时新请求抢走资源。

  • 指数退避:wait_exponential 让等待时间逐渐拉长,避免集体踩踏。

  • 有限重试 + reraise:防止无限重试,失败最终向上抛出,由业务层决定降级或放弃。


🧩 5. 完整并发控制方案流程图(逻辑链)

[计算semaphore size]
[发起请求] → 获取信号量(无空位则等待)
[执行 HTTP 调用] → 成功 → 释放信号量 → 返回
      ↓ 429?
[tenacity 判断重试] → 等待指数退避
[重新执行 HTTP 调用](仍在信号量内)
      ... 直到成功或达到最大重试次数

💡 6. 生产环境里我踩过的坑 & 最佳实践

  • 动态调整并发数:请求耗时 D 和 token 消耗 T 不是固定的。我会在代码里定时计算最近 1 分钟的平均耗时和平均 token,然后周期性更新 semaphore._value(或重建 Semaphore)。但要注意并发安全,最好通过事件循环安全的方式。

  • 别只看平均值:P99 延迟突然变大会瞬间超额,安全系数 0.8 就是为这准备的。

  • 加上监控告警:记录每分钟实际请求数、token 消耗、429 次数、重试成功率。如果 429 频繁,说明并发数还是太激进,得再缩一点。

  • Retry-After 头:如果 API 返回了 Retry-After,可以写一个自定义 wait 函数,优先用这个值,比盲目退避更精准。


最后留一句很实在的话:限额策略永远是“估算+监控+迭代”出来的,别指望一次算准管一辈子。我通常先按公式算个安全值上线,然后根据线上 429 比例慢慢把并发数往上调,直到找到一个不爆炸的甜点。 如果面试官想听,我还能聊聊如何用令牌桶在本地再做一层整形,把突发请求拍得平滑一点。