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

Binance多线程Socket场景下analyze函数无法并发调用问题咨询

问题原因

当前代码无法并发执行analyze的核心问题有两个:

  1. analyze是同步函数,内部使用的time.sleep(5)是阻塞调用,会直接卡住整个asyncio事件循环,所有异步任务都会暂停直到sleep结束
  2. 收到websocket推送数据后,同步调用analyze,必须等该函数执行完成后才会继续接收下一条推送,自然无法并行处理两个交易对的数据
    运行效果截图

修复方案

按照如下方式调整代码即可实现analyze独立并发调用:

  1. 将analyze改造为异步函数,把阻塞的time.sleep替换为非阻塞的await asyncio.sleep
  2. 收到数据后将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 22:24:03