如何在异步协程中添加延时以规避Binance API请求权重限制
问题场景
并发拉取Binance加密货币交易对历史K线时触发API权重限制,IP被封禁,报错如下:
APIError(code=-1003): Way too much request weight used; IP banned until 1629399758399. Please use the websocket for live updates to avoid bans.
Binance公开API请求权重限制为每分钟1200,当前实现为全量交易对一次性并发请求,导致瞬间权重打满被封。
修复方案
核心思路是平滑请求节奏,限制单位时间内的请求总量,具体修改点如下:
- 复用全局AsyncClient实例,避免每个协程重复创建客户端造成的额外开销
- 用
asyncio.Semaphore限制同时执行的请求协程数量,避免瞬时并发过高 - 分批处理交易对,每处理完一批主动添加延时,将每分钟请求总量控制在1200以内
- 增加异常捕获逻辑,避免单个请求失败导致整个任务终止
修改后可运行代码
import pandas as pd import asyncio import asyncpg import time from binance.client import AsyncClient import config # 限流配置参数 MAX_CONCURRENT = 200 # 最大同时并发请求数 BATCH_SIZE = 200 # 每批处理的交易对数量 BATCH_WAIT = 10 # 每批处理完后的等待时间(秒) # 全局信号量控制并发 semaphore = asyncio.Semaphore(MAX_CONCURRENT) # 全局复用Binance客户端,避免重复创建连接 binance_client = None async def main(): global binance_client # 初始化全局Binance客户端 binance_client = await AsyncClient.create(config.BINANCE_API_KEY, config.BINANCE_SECRET_KEY) # 创建数据库连接池 pool = await asyncpg.create_pool( user=config.DB_USER, password=config.DB_PASS, database=config.DB_NAME, host=config.DB_HOST, command_timeout=60 ) # 读取数据库中存储的交易对列表 async with pool.acquire() as connection: cryptos = await connection.fetch("SELECT * FROM crypto") symbols = {crypto['id']: crypto['symbol'] for crypto in cryptos} symbol_list = list(symbols.items()) total = len(symbol_list) # 分批处理交易对 for i in range(0, total, BATCH_SIZE): batch = symbol_list[i:i+BATCH_SIZE] print(f"Processing batch {i//BATCH_SIZE + 1}, {len(batch)} symbols") tasks = [asyncio.create_task(get_price(pool, crypto_id, symbol)) for crypto_id, symbol in batch] await asyncio.gather(*tasks, return_exceptions=True) # 非最后一批则主动等待,避免请求过密 if i + BATCH_SIZE < total: await asyncio.sleep(BATCH_WAIT) await binance_client.close_connection() print(f"Finalized all. Retrieved price data of {total} outputs.") async def get_price(pool, crypto_id, symbol): async with semaphore: try: candlesticks = [] async for kline in await binance_client.get_historical_klines_generator( symbol, AsyncClient.KLINE_INTERVAL_1HOUR, "18 Aug, 2021", "19 Aug, 2021" ): candlesticks.append(kline) df = pd.DataFrame( candlesticks, columns = ["date","open","high","low","close","volume","Close time","Quote Asset Volume","Number of Trades","Taker buy base asset volume","Taker buy quote asset volume","Ignore"] ) df["date"] = pd.to_datetime(df.loc[:, "date"], unit ='ms') df.drop(columns=['Close time','Ignore', 'Quote Asset Volume', 'Number of Trades', 'Taker buy base asset volume', 'Taker buy quote asset volume'], inplace=True) df.loc[:, "id"] = crypto_id # 此处可添加df写入数据库的逻辑 print(f"Successfully fetched {symbol} data, {len(df)} records") except Exception as e: print(f"Unable to get {symbol} prices due to {e.__class__}: {e}") if __name__ == "__main__": start = time.time() asyncio.run(main()) end = time.time() print(f"Took {end - start} seconds.")
参数调整说明
如果运行后仍触发限流,可调小MAX_CONCURRENT和BATCH_SIZE,增大BATCH_WAIT;也可以读取Binance接口响应头的x-mbx-used-weight-1m字段动态判断当前已用权重,达到1000左右时主动等待到下一分钟再继续请求,适配度更高。
内容的提问来源于stack exchange,提问作者JC Martinez
相关产品推荐
相关产品推荐

