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

如何在异步协程中添加延时以规避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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.07 13:30:03