必须将令牌桶算法深度嵌入协程调度链路,使每个await前显式申请令牌:用asyncio.Lock+asyncio.sleep实现原子填充,禁用同步操作,多进程需Redis Lua保证一致性。

要对高频并发的 async 任务做硬性速率限制,不能只靠“等一等”或“数个数”,必须把令牌桶算法深度嵌入协程调度链路中——核心是让每个 await 操作在真正执行前,先向令牌桶申请许可,无令牌则挂起,不阻塞事件循环。
令牌桶需与事件循环协同工作
Python 中不能用普通线程安全计数器(如 threading.Lock + int)实现令牌桶,因为 async 任务不共享线程上下文。必须使用 asyncio 兼容的原子操作原语:
- 用 asyncio.Semaphore 管理并发槽位数(对应桶容量),但它只控制“同时运行几个”,不解决“单位时间放多少令牌”问题
- 真正可行的是基于 asyncio.Queue 或自定义 TokenBucket 类,内部用 asyncio.Event + asyncio.create_task 启动后台填充协程,按固定间隔(如每秒 10 个)向队列注入令牌
- 关键细节:填充协程必须用 asyncio.sleep() 而非 time.sleep(),否则会阻塞整个事件循环
await 前必须显式调用 acquire_token()
不能把限流逻辑藏在装饰器或中间件里“自动生效”,必须让每个高风险 async 任务显式 await 一个令牌获取操作。示例结构如下:
class TokenBucket:
def __init__(self, capacity: int, refill_rate: float):
self.capacity = capacity
self.tokens = capacity
self.refill_rate = refill_rate
self.last_refill = asyncio.get_event_loop().time()
self._lock = asyncio.Lock()
<pre class="brush:php;toolbar:false;">async def acquire(self):
while True:
async with self._lock:
now = asyncio.get_event_loop().time()
elapsed = now - self.last_refill
new_tokens = int(elapsed * self.refill_rate)
if new_tokens > 0:
self.tokens = min(self.capacity, self.tokens + new_tokens)
self.last_refill = now
if self.tokens > 0:
self.tokens -= 1
return True
# 无令牌时,等待 10ms 后重试,避免忙等
await asyncio.sleep(0.01)
调用方写法必须是:await bucket.acquire() —— 这个 await 就是硬性管控点,它不返回就卡住当前协程,但不阻塞其他任务。
避免常见陷阱:同步操作污染事件循环
很多团队失败是因为在 acquire() 内部混入了同步耗时操作,例如:
- 调用 redis.incr()(未用 aioredis)→ 实际是同步阻塞调用 → 整个事件循环卡死
- 用 datetime.now() 替代 asyncio.get_event_loop().time() → 时间精度低,且无实际危害但暴露设计粗糙
- 令牌桶状态存在本地内存 → 多进程部署时失效 → 必须用 Redis 的 Lua 脚本保证原子性(EVAL "local tokens = ..." 2 rate limit_key)
与业务逻辑解耦但强绑定
令牌桶本身不关心任务内容,但它的 acquire() 必须成为任务入口的强制前置步骤。推荐模式:
- 定义统一的 async 工厂函数:async def run_with_rate_limit(task_coro, bucket)
- 内部先 await bucket.acquire(),再 await task_coro(),异常时也确保令牌归还(可选)
- 禁止直接 await 业务协程,所有高频调用路径必须走该工厂函数
- FastAPI 中可封装为依赖项(Depends),让路由函数签名强制携带 bucket 实例
不复杂但容易忽略:硬性管控不是加个装饰器就能生效,而是把令牌获取变成 await 链上不可跳过的环节。

















