如何避免Python异步协程阻塞?实现websocket客户端后台运行
解决WebSocket客户端阻塞主线程的问题
你的问题核心在于feed.__init__里调用了asyncio.get_event_loop().run_until_complete(self.socket())——因为self.socket()是一个无限循环的协程(一直在接收WebSocket消息),run_until_complete会一直阻塞到协程结束,自然永远到不了后面的print("hello")循环。
下面给你两种可行的解决方案,你可以根据自己的代码架构选择:
方案一:全异步架构(推荐)
把整个主逻辑改成异步,利用asyncio的任务并发能力,让WebSocket客户端和打印逻辑同时运行:
import websockets import asyncio import json class feed: def __init__(self, symbol): self.symbol = symbol self.uri = f"wss://www.bitmex.com/realtime?subscribe=instrument,quote:{symbol}" self.last_message = None async def socket(self): async with websockets.connect(self.uri) as websocket: while True: msg = await websocket.recv() self.process_msg(json.loads(msg)) def process_msg(self, msg): try: data = msg['data'][0] if data['symbol'] == self.symbol: self.last_message = msg # 可选:在这里处理消息,比如打印更新 # print(f"Received update for {self.symbol}: {data}") except Exception as e: print(f"Unhandled error processing message: {e}") async def print_hello_loop(): while True: print("hello") await asyncio.sleep(1) # 避免无限制打印占用CPU async def main(): myFeed = feed('XBTUSD') # 将WebSocket任务提交到事件循环后台运行 ws_task = asyncio.create_task(myFeed.socket()) # 同时运行WebSocket任务和打印任务 await asyncio.gather(ws_task, print_hello_loop()) if __name__ == "__main__": asyncio.run(main())
为什么这样可行?
asyncio.create_task会把WebSocket协程包装成一个后台任务,交给事件循环调度,不会阻塞当前代码。asyncio.gather会同时等待多个协程完成,这样两个任务就能并发执行了。- 给
print_hello_loop加await asyncio.sleep(1)是为了避免疯狂打印占用过多系统资源,你可以根据需求调整间隔。
方案二:后台线程运行事件循环(适合不想改全异步的场景)
如果你希望保持main函数是同步的,可以把WebSocket的事件循环放到后台线程中运行,这样主线程就能继续执行打印逻辑:
import websockets import asyncio import json import threading import time class feed: def __init__(self, symbol): self.symbol = symbol self.uri = f"wss://www.bitmex.com/realtime?subscribe=instrument,quote:{symbol}" self.last_message = None # 创建独立的事件循环并在后台线程运行 self.loop = asyncio.new_event_loop() # 设置daemon=True,确保主线程退出时后台线程也会终止 threading.Thread(target=self._run_event_loop, daemon=True).start() # 将WebSocket协程提交到后台循环 asyncio.run_coroutine_threadsafe(self.socket(), self.loop) def _run_event_loop(self): asyncio.set_event_loop(self.loop) self.loop.run_forever() async def socket(self): async with websockets.connect(self.uri) as websocket: while True: msg = await websocket.recv() self.process_msg(json.loads(msg)) def process_msg(self, msg): try: data = msg['data'][0] if data['symbol'] == self.symbol: self.last_message = msg except Exception as e: print(f"Unhandled error processing message: {e}") def main(): myFeed = feed('XBTUSD') while True: print("hello") time.sleep(1) # 避免CPU占用过高 if __name__ == "__main__": main()
注意事项:
- 后台线程用
daemon=True,这样当主线程结束时,后台线程会自动退出,不会让程序挂着。 - 如果需要在主线程和后台线程之间共享数据(比如
self.last_message),要注意线程安全,可以用threading.Lock来保护共享变量,避免竞态条件。
额外建议
- 把原来的
except:改成except Exception as e:,这样能捕获并打印具体的异常信息,方便调试问题,而不是只输出模糊的"Unhandled"。 - 如果你需要处理WebSocket的断开重连,可以在
socket协程里加异常捕获,断开后自动重新连接。
内容的提问来源于stack exchange,提问作者cjm2671
相关产品推荐
相关产品推荐

