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

如何结合aiohttp实现Riot API的多维度速率限制?

解决Riot API双维度限流的异步实现方案

针对你遇到的异步环境下Riot API双维度限流(每秒20次、120秒内100次)问题,单纯用信号量无法满足多窗口的限流要求,需要实现一个双维度的主动限流器,结合兜底的429重试逻辑,同时优化数据库写入效率。

核心思路

Riot的双维度限流需要同时满足两个条件:

  • 任意1秒内请求数不超过20
  • 任意120秒内请求数不超过100

主动限流的关键是维护两个时间窗口的请求记录,每次发起请求前先校验两个窗口的剩余额度,不足则等待到可用时间,避免触发429导致账号被拉黑。

实现双维度限流器

用asyncio的锁保护请求计数器,结合双队列分别维护两个时间窗口的请求时间戳:

import asyncio
from collections import deque
from datetime import datetime, timedelta

class DualRateLimiter:
    def __init__(self, qps_limit: int, window_limit: int, window_seconds: int):
        self.qps_limit = qps_limit
        self.window_limit = window_limit
        self.window_seconds = window_seconds
        # 存储最近1秒的请求时间戳
        self.qps_queue = deque()
        # 存储最近window_seconds秒的请求时间戳
        self.window_queue = deque()
        self.lock = asyncio.Lock()

    async def acquire(self):
        async with self.lock:
            now = datetime.now()
            # 清理过期的请求记录
            self._clean_expired(now)

            # 计算需要等待的时间,取两个限制中最长的等待时长
            wait_time = self._calculate_wait_time(now)

            if wait_time > 0:
                await asyncio.sleep(wait_time)
                # 等待后重新清理过期记录,避免期间有其他请求插入
                now = datetime.now()
                self._clean_expired(now)

            # 记录当前请求时间
            self.qps_queue.append(now)
            self.window_queue.append(now)

    def _clean_expired(self, now: datetime):
        # 清理1秒外的QPS记录
        while self.qps_queue and now - self.qps_queue[0] > timedelta(seconds=1):
            self.qps_queue.popleft()
        # 清理时间窗口外的记录
        while self.window_queue and now - self.window_queue[0] > timedelta(seconds=self.window_seconds):
            self.window_queue.popleft()

    def _calculate_wait_time(self, now: datetime) -> float:
        wait_time = 0
        # 处理QPS限制
        if len(self.qps_queue) >= self.qps_limit:
            next_qps_time = self.qps_queue[0] + timedelta(seconds=1)
            wait_time = max(wait_time, (next_qps_time - now).total_seconds())
        # 处理时间窗口限制
        if len(self.window_queue) >= self.window_limit:
            next_window_time = self.window_queue[0] + timedelta(seconds=self.window_seconds)
            wait_time = max(wait_time, (next_window_time - now).total_seconds())
        return wait_time

结合aiohttp封装请求

把限流器和aiohttp结合,同时保留429的兜底重试逻辑(防止主动限流有误差):

import aiohttp

class RiotAPIClient:
    def __init__(self, api_key: str):
        self.api_key = api_key
        # 初始化对应Riot API的限流参数
        self.rate_limiter = DualRateLimiter(qps_limit=20, window_limit=100, window_seconds=120)
        self.session = aiohttp.ClientSession(headers={"X-Riot-Token": api_key})

    async def fetch(self, url: str):
        while True:
            # 先获取限流许可
            await self.rate_limiter.acquire()
            try:
                async with self.session.get(url) as response:
                    if response.status == 200:
                        return await response.json()
                    elif response.status == 429:
                        # 按API返回的重试时间等待,加1秒避免刚好卡点
                        retry_after = int(response.headers.get("Retry-After", 1))
                        await asyncio.sleep(retry_after + 1)
                    else:
                        print(f"请求失败:状态码{response.status},URL: {url}")
                        return None
            except aiohttp.ClientError as e:
                print(f"请求异常:{str(e)}")
                await asyncio.sleep(1)

    async def close(self):
        await self.session.close()

适配你的业务流程

注意数据库操作要用异步驱动(比如aiomysql),避免阻塞事件循环;同时可以批量积累数据后再插入,提升写入效率:

import aiomysql

# 异步数据库存储示例
async def save_to_db(pool: aiomysql.Pool, match_data, player_data):
    async with pool.acquire() as conn:
        async with conn.cursor() as cur:
            # 这里写你的插入逻辑,建议批量插入
            await cur.execute("INSERT INTO matches (...) VALUES (...)", (match_data["id"], ...))
            await cur.execute("INSERT INTO players (...) VALUES (...)", (player_data["puuid"], ...))
            await conn.commit()

async def process_match(match_id: str, client: RiotAPIClient, db_pool: aiomysql.Pool):
    # 抓取匹配数据
    match_data = await client.fetch(f"https://api.riotgames.com/lol/match/v5/matches/{match_id}")
    if not match_data:
        return
    # 批量处理玩家信息
    player_tasks = []
    for participant in match_data["info"]["participants"]:
        puuid = participant["puuid"]
        player_tasks.append(client.fetch(f"https://api.riotgames.com/lol/summoner/v4/summoners/by-puuid/{puuid}"))
    player_datas = await asyncio.gather(*player_tasks)
    # 批量存入数据库
    for player_data in player_datas:
        if player_data:
            await save_to_db(db_pool, match_data, player_data)

async def main():
    api_key = "你的API密钥"
    client = RiotAPIClient(api_key)
    # 初始化异步数据库连接池
    db_pool = await aiomysql.create_pool(
        host="localhost",
        port=3306,
        user="your_user",
        password="your_pwd",
        db="your_db"
    )
    # 同步获取匹配列表(可以改成异步请求)
    match_list = get_match_list_sync()  # 替换成你的同步获取逻辑
    # 并发处理匹配
    tasks = [process_match(match_id, client, db_pool) for match_id in match_list]
    await asyncio.gather(*tasks)
    # 关闭资源
    await client.close()
    db_pool.close()
    await db_pool.wait_closed()

if __name__ == "__main__":
    asyncio.run(main())

额外优化建议

  1. 批量写入数据库:不要每条数据单独插入,积累一定数量后批量提交,减少数据库连接开销
  2. 日志系统:添加详细的日志记录,包括限流等待时间、请求状态、数据库操作结果,方便排查问题
  3. 动态调整限流参数:如果API密钥的限额变化,可以把限流参数做成可配置的,不用硬编码

内容的提问来源于stack exchange,提问作者Gustavo Feijó

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 16:53:11