如何结合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())
额外优化建议
- 批量写入数据库:不要每条数据单独插入,积累一定数量后批量提交,减少数据库连接开销
- 日志系统:添加详细的日志记录,包括限流等待时间、请求状态、数据库操作结果,方便排查问题
- 动态调整限流参数:如果API密钥的限额变化,可以把限流参数做成可配置的,不用硬编码
内容的提问来源于stack exchange,提问作者Gustavo Feijó
相关产品推荐
相关产品推荐

