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

如何避免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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.09 06:17:53