使用python-binance调用API时请求过快触发队列溢出错误如何解决
问题原因
你使用的是基于async/await的异步代码,time.sleep()是同步阻塞函数,在异步协程中调用时会直接卡住整个事件循环,无法实现预期的速率控制效果。
另外你的当前逻辑是持续接收WebSocket服务端推送的行情数据,原生推送频率本身就远高于你需要的每秒1次,直接加休眠反而会导致客户端缓冲区消息堆积,后续还是会高频处理触发API队列溢出限制。
解决方案
推荐两种适配不同需求的实现方式:
方式1:跳过冗余消息,仅每秒处理1条(适合不需要保留全量推送数据的场景)
import asyncio import time # 初始化上次处理时间 last_process_time = 0 while True: await socket.__aenter__() msg = await socket.recv() current_time = time.time() # 距离上次处理不足1秒直接跳过当前消息 if current_time - last_process_time < 1: continue frame = createframe(msg) frame.to_sql(symbol, engine, if_exists="append", index=False) print(frame) # 更新上次处理时间 last_process_time = current_time
方式2:本地缓存消息,每秒批量写入(适合需要保留全量推送数据的场景)
import asyncio from collections import deque # 初始化本地消息缓存队列 msg_queue = deque() # 单独的批量写库协程 async def write_to_db(): while True: if msg_queue: # 一次性取出所有缓存的消息批量处理 frames = [createframe(msg) for msg in msg_queue] # 批量写入数据库,比单条写入性能更高 for frame in frames: frame.to_sql(symbol, engine, if_exists="append", index=False) print(frame) # 清空缓存队列 msg_queue.clear() # 异步休眠1秒再处理下一批,不阻塞消息接收 await asyncio.sleep(1) # 原消息接收协程 async def recv_msg(): while True: await socket.__aenter__() msg = await socket.recv() # 消息先存入缓存队列 msg_queue.append(msg) # 事件循环中同时运行两个协程 loop = asyncio.get_event_loop() loop.run_until_complete(asyncio.gather(recv_msg(), write_to_db()))
内容的提问来源于stack exchange,提问作者Donnerhorn
相关产品推荐
相关产品推荐

