Binance多线程Socket场景下analyze函数无法并发调用问题咨询
问题原因
当前代码无法并发执行analyze的核心问题有两个:
analyze是同步函数,内部使用的time.sleep(5)是阻塞调用,会直接卡住整个asyncio事件循环,所有异步任务都会暂停直到sleep结束- 收到websocket推送数据后,同步调用
analyze,必须等该函数执行完成后才会继续接收下一条推送,自然无法并行处理两个交易对的数据
修复方案
按照如下方式调整代码即可实现analyze独立并发调用:
- 将
analyze改造为异步函数,把阻塞的time.sleep替换为非阻塞的await asyncio.sleep - 收到数据后将
analyze封装为独立asyncio任务提交,无需等待执行完成,主接收循环可以立刻处理下一条数据
修改后完整代码
import asyncio from binance import AsyncClient, BinanceSocketManager from datetime import datetime async def analyze(res): try: kline = res['k'] if kline['x']: # 蜡烛图已完成 print('{} start_sleeping {} {}'.format( datetime.now(), kline['s'], datetime.fromtimestamp(kline['t'] / 1000), )) # 替换为非阻塞异步休眠,会让出事件循环给其他任务执行 await asyncio.sleep(5) print('{} finish_sleeping {}'.format(datetime.now(), kline['s'])) except Exception as e: print(f"处理交易对数据出错: {str(e)}") async def open_binance_stream(symbol): client = await AsyncClient.create() bm = BinanceSocketManager(client) ts = bm.kline_socket(symbol) async with ts as tscm: while True: res = await tscm.recv() # 创建独立任务运行analyze,不阻塞主接收逻辑 asyncio.create_task(analyze(res)) await client.close_connection() async def main(): t1 = asyncio.create_task(open_binance_stream('ETHBTC')) t2 = asyncio.create_task(open_binance_stream('XRPBTC')) await asyncio.gather(t1, t2) if __name__ == "__main__": asyncio.run(main())
补充说明
如果后续analyze需要加入CPU密集型计算逻辑,建议使用asyncio.to_thread()将计算部分丢到独立线程运行,避免阻塞异步事件循环。
内容的提问来源于stack exchange,提问作者stan
相关产品推荐
相关产品推荐

