Python使用websockets与asyncio监听多Socket失效问题及解决方案咨询
解决多Websocket监听问题的方案
嘿,我来帮你搞定这个问题!你遇到的“只能监听第一个Socket”的情况,大概率是因为你在处理第一个Socket的无限循环时,没有给其他异步任务执行的机会。Asyncio是协作式多任务,只要你的代码里有await(比如websocket的接收操作),就不会完全阻塞,但如果你的任务启动方式不对,还是会出现“卡住第一个”的情况。下面给你两种可行的解决办法:
优先方案:用Asyncio原生方式管理多任务
这是最符合Asyncio设计理念的方案,不需要额外引入线程,性能也更好。核心思路是:把所有Websocket的订阅任务都创建好,然后用asyncio.gather()让它们同时运行。
修正后的代码示例
import asyncio import json import websockets class SocketManager: def __init__(self): self.tasks = [] async def subscribe(self, event): # 替换成你的实际Websocket地址 ws_uri = "ws://your-target-websocket-server" async with websockets.connect(ws_uri) as websocket: # 发送订阅请求 payload = json.dumps(event) await websocket.send(payload) # 无限循环监听消息——这里的await是关键,会让出CPU给其他任务 while True: message = await websocket.recv() print(f"[{event['channel']}] 收到消息: {message}") def add_socket_task(self, event): # 只创建任务,不立即执行,避免阻塞后续任务添加 current_loop = asyncio.get_event_loop() self.tasks.append(current_loop.create_task(self.subscribe(event))) async def run_all_sockets(self): # 同时运行所有已添加的Socket任务 await asyncio.gather(*self.tasks) # 用法示例 if __name__ == "__main__": manager = SocketManager() # 添加多个Socket监听任务 manager.add_socket_task({"channel": "trade_btc"}) manager.add_socket_task({"channel": "trade_eth"}) manager.add_socket_task({"channel": "order_book"}) # 启动所有任务 asyncio.run(manager.run_all_sockets())
为什么这个方案能解决问题?
- 你之前的代码可能是在调用
subscribe时直接await了,导致第一个无限循环占用了整个事件循环,后面的任务根本没机会启动。 - 现在我们先把所有任务都创建好存入列表,最后用
asyncio.gather()一次性启动所有任务。每个subscribe里的await websocket.recv()会主动让出CPU,让其他任务也能执行。
备选方案:为每个Socket单独使用线程
如果你的场景中存在无法异步化的同步阻塞代码(比如某些老库不支持Asyncio),可以给每个Socket单独开一个线程,每个线程运行自己的Asyncio事件循环。
线程方案的代码示例
import asyncio import json import websockets import threading def run_single_socket(event): async def subscribe(): ws_uri = "ws://your-target-websocket-server" async with websockets.connect(ws_uri) as websocket: payload = json.dumps(event) await websocket.send(payload) while True: message = await websocket.recv() print(f"[{event['channel']}] 收到消息: {message}") # 每个线程创建独立的事件循环 loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) loop.run_until_complete(subscribe()) # 用法示例 if __name__ == "__main__": # 为每个Socket创建线程 thread1 = threading.Thread(target=run_single_socket, args=({"channel": "trade_btc"},)) thread2 = threading.Thread(target=run_single_socket, args=({"channel": "trade_eth"},)) thread1.start() thread2.start() # 等待所有线程结束(可选,根据你的需求) thread1.join() thread2.join()
注意事项
- 线程间的资源共享需要注意加锁(比如如果多个线程要操作同一个全局变量),但如果只是各自监听独立的Socket,基本不需要额外处理。
- 这种方案的资源开销比纯Asyncio大,所以如果没有特殊需求,优先用第一种方案。
内容的提问来源于stack exchange,提问作者joseRo
相关产品推荐
相关产品推荐

