You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何修复基于asyncio的限流API请求代码停滞问题?

问题分析与修复方案

问题根源

  1. 全局变量竞争:多个异步任务无保护地修改active_calls、completed_calls等全局变量,导致状态计算混乱——比如多个任务同时通过槽位检查,瞬间占用超过上限的资源。
  2. 时间窗口逻辑错误:timeout_start初始已赋值,任务完成时的if timeout_start is None判断永远不会触发,导致时间窗口无法更新;重置completed_calls时未同步更新时间窗口起始点,后续等待时间计算完全错误。
  3. 休眠唤醒后的抢占问题:时间窗口到期后,所有等待的任务同时被唤醒并抢占槽位,导致资源占用失控,进而引发后续任务的无限等待。

修复方案

使用异步安全的令牌桶算法实现限流,通过锁保护状态操作,确保逻辑正确。同时复用HTTP会话提升性能。

修复后的代码

import asyncio
import aiohttp
from datetime import datetime, timedelta

class TokenBucket:
    def __init__(self, max_tokens, refill_interval):
        self.max_tokens = max_tokens
        self.refill_interval = refill_interval
        self.tokens = max_tokens
        self.last_refill = datetime.now()
        self.lock = asyncio.Lock()

    async def acquire(self):
        async with self.lock:
            now = datetime.now()
            time_since_refill = now - self.last_refill
            
            # 到达时间窗口,重置令牌和起始时间
            if time_since_refill >= self.refill_interval:
                self.tokens = self.max_tokens
                self.last_refill = now

            # 无令牌可用时,等待到下一次窗口刷新
            while self.tokens <= 0:
                wait_time = self.refill_interval - time_since_refill
                await asyncio.sleep(wait_time.total_seconds())
                # 重新计算时间差和令牌数
                now = datetime.now()
                time_since_refill = now - self.last_refill
                if time_since_refill >= self.refill_interval:
                    self.tokens = self.max_tokens
                    self.last_refill = now

            self.tokens -= 1

# 限流配置:每3秒最多5次请求
TIME_WINDOW = timedelta(seconds=3)
MAX_CALLS = 5

async def task(i, bucket, session):
    await bucket.acquire()
    print(f'Task {i} start - 剩余令牌: {bucket.tokens}')

    try:
        async with session.get('https://www.google.com/') as response:
            print(f'Task {i} 完成 - 状态码: {response.status}')
    finally:
        # 令牌按时间窗口自动补充,无需手动释放
        pass

async def main():
    bucket = TokenBucket(MAX_CALLS, TIME_WINDOW)
    # 复用ClientSession,避免重复创建连接
    async with aiohttp.ClientSession() as session:
        await asyncio.gather(*[task(i, bucket, session) for i in range(20)])

asyncio.run(main())

关键修复点

  • 异步安全的状态管理:用asyncio.Lock确保同一时间只有一个任务能修改令牌状态,彻底解决竞争问题。
  • 正确的令牌补充逻辑:每次获取令牌时先检查时间窗口,自动重置令牌数,无需手动维护完成计数。
  • 复用HTTP会话:所有任务共享一个ClientSession,减少连接开销,提升请求效率。
  • 清晰的状态封装:将限流逻辑封装到TokenBucket类中,代码更易维护和扩展。

内容的提问来源于stack exchange,提问作者Gregor

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.24 17:37:15