如何利用多asyncio协程向同一Queue写入WebSocket实时数据?
解决方案
核心问题分析
你当前的代码存在两个关键问题:
asyncio.run()是阻塞式调用,每次只能启动一个协程并等待它执行完毕。但你的data_stream是无限循环的WebSocket数据监听逻辑,后续的B_Orderbook和C_Orderbook永远不会被执行。- 使用了线程安全的
queue.Queue,但在异步协程场景下,应该用专为asyncio设计的asyncio.Queue——它支持await式的读写操作,不会阻塞事件循环。
修正后的实现代码
import asyncio async def main(): # 替换为asyncio专属队列,适配异步场景 global_queue = asyncio.Queue() # 同时启动三个数据监听协程 await asyncio.gather( A_Orderbook.data_stream(global_queue), B_Orderbook.data_stream(global_queue), C_Orderbook.data_stream(global_queue), return_exceptions=True # 可选:单个协程出错时不终止其他协程 ) if __name__ == '__main__': try: asyncio.run(main()) except KeyboardInterrupt: print("程序已终止")
对data_stream方法的要求
确保你的A_Orderbook.data_stream等方法是异步函数,并且写入队列时使用await queue.put(data),示例如下:
import aiohttp class A_Orderbook: @staticmethod async def data_stream(queue): # WebSocket场景示例 async with aiohttp.ClientSession() as session: async with session.ws_connect("wss://a.example.com/ws") as ws: async for msg in ws: if msg.type == aiohttp.WSMsgType.TEXT: data = msg.json() await queue.put(data) elif msg.type == aiohttp.WSMsgType.ERROR: break # 异步轮询场景示例 @staticmethod async def data_stream(queue): while True: async with aiohttp.ClientSession() as session: async with session.get("https://a.example.com/data") as resp: data = await resp.json() await queue.put(data) await asyncio.sleep(1) # 异步休眠,避免高频请求
关键说明
asyncio.gather():可以同时运行多个协程,等待所有协程完成(如果是无限循环逻辑则会持续运行)。return_exceptions=True:如果某个数据源连接中断抛出异常,其他协程仍能继续运行,便于排查单个数据源的问题。asyncio.Queue的读写必须用await queue.put()和await queue.get(),不能直接使用同步的put()/get(),否则会阻塞事件循环。
内容的提问来源于stack exchange,提问作者Nathan
相关产品推荐
相关产品推荐

