Asyncio同时维护两个WebSocket连接时无法执行第三个异步函数
问题分析与解决方案
核心错误点
- 事件循环被阻塞:如果你的WebSocket连接用了同步式死循环(比如
while True持续收消息但不让出事件循环),会直接霸占异步事件循环,导致checking()完全没机会执行。异步编程的核心是任务间要主动让出执行权,否则单个任务会占用所有资源。 checking()未被正确注册到事件循环:你可能只定义了checking()函数,但没通过asyncio.create_task()这类方法把它加入异步事件循环;或者启动顺序错误,导致它根本没启动。- 全局变量的异步适配问题:直接用全局变量传递数据时,
checking()如果没有合理的触发逻辑(比如异步通知而非定时轮询),要么拿不到新数据,要么轮询逻辑本身又阻塞了事件循环。
正确实现思路
放弃全局变量,改用异步队列(asyncio.Queue)传递数据,同时确保所有异步任务都被正确注册到事件循环并行运行:
- 两个WebSocket任务负责接收数据,将数据放入队列;
checking()任务持续从队列取数据处理;- 用
asyncio.create_task()启动所有任务,通过asyncio.gather()维持任务运行。
示例代码
import asyncio import websockets async def kraken_websocket(queue): # 连接Kraken WebSocket uri = "wss://ws.kraken.com" async with websockets.connect(uri) as ws: # 订阅BTC/USDT行情(示例订阅指令) await ws.send('{"event":"subscribe","pair":["XBT/USD"],"subscription":{"name":"ticker"}}') # 异步接收消息,自动让出事件循环 async for msg in ws: await queue.put(("kraken", msg)) async def binance_websocket(queue): # 连接Binance WebSocket uri = "wss://stream.binance.com:9443/ws/btcusdt@ticker" async with websockets.connect(uri) as ws: async for msg in ws: await queue.put(("binance", msg)) async def checking(queue): while True: # 异步等待队列新数据,不会阻塞事件循环 source, data = await queue.get() # 这里写你的数据处理逻辑 print(f"从{source}收到数据:{data[:60]}...") # 标记任务完成(可选,用于队列join()方法) queue.task_done() async def main(): # 创建异步队列用于数据传递 data_queue = asyncio.Queue() # 启动三个异步任务 kraken_task = asyncio.create_task(kraken_websocket(data_queue)) binance_task = asyncio.create_task(binance_websocket(data_queue)) checking_task = asyncio.create_task(checking(data_queue)) # 等待所有任务持续运行(手动终止前不会退出) await asyncio.gather(kraken_task, binance_task, checking_task) if __name__ == "__main__": asyncio.run(main())
关键说明
async for处理WebSocket消息:这种方式会在每次接收消息后主动让出事件循环,让checking()有执行机会;- 异步队列替代全局变量:队列是异步安全的,能避免全局变量的竞态问题,数据传递更可靠;
- 任务并行启动:通过
create_task()将三个任务都加入事件循环,由asyncio自动调度并行运行。
内容的提问来源于stack exchange,提问作者Treasure Hunter
相关产品推荐
相关产品推荐

