如何立即终止所有异步运行任务?附Python WebSocket代码问题
问题分析
你的代码核心问题有两个:
await websocket.recv()是无阻塞上限的调用,Bybit的WebSocket会先返回订阅确认消息(而非目标数据),后续还可能发送心跳包,这些都会被你的代码接收但不会计入ws_responses,导致循环一直等待,直到真正的深度数据到来(可能延迟40秒)。- 没有超时机制,一旦WebSocket连接异常或迟迟不返回目标数据,程序会无限挂起。
解决方案
给接收消息的逻辑加上超时限制,同时精准过滤目标消息,确保收集到所需数据后立即终止连接并返回:
import asyncio import json import websockets import traceback async def subscribe(url, subs, timeout=10): ws_responses = [] target_count = len(subs) async with websockets.connect(url) as websocket: # 发送所有订阅请求 for sub in subs: sub_str = json.dumps(sub) await websocket.send(sub_str) while len(ws_responses) < target_count: try: # 给recv加上超时,避免无限等待 rsp_str = await asyncio.wait_for(websocket.recv(), timeout=timeout) rsp = json.loads(rsp_str) except asyncio.TimeoutError: # 超时未收到目标数据,抛出异常让上层重连 raise TimeoutError(f"Timeout after {timeout}s waiting for target data") # 精准过滤目标数据:Bybit的depth订阅返回的消息中,topic字段会匹配订阅的topic if rsp.get('topic') == subs[0]['topic'] and 'data' in rsp: ws_responses.append(rsp['data']['s']) # 收集到足够数据后,主动关闭连接并跳出循环 await websocket.close() break return ws_responses def main(): market_url = 'wss://stream.bybit.com/spot/quote/ws/v2' market_subs = [{ "topic": "depth", "event": "sub", "params": { "symbol": "BTCUSDT", }} ] while True: try: response = asyncio.get_event_loop().run_until_complete( subscribe(market_url, market_subs)) return response except Exception as e: traceback.print_exc() print('websocket connection error. reconnect rightnow') if __name__ == "__main__": print(main())
关键改动说明
- 用
asyncio.wait_for()给websocket.recv()设置超时(默认10秒),超时后抛出异常触发重连,避免无限等待。 - 消息过滤逻辑优化:不仅检查
data字段,还匹配topic字段,确保只处理订阅的深度数据,忽略订阅确认和心跳包。 - 收集到目标数据后,主动调用
websocket.close()关闭连接,确保异步上下文管理器快速退出。
内容的提问来源于stack exchange,提问作者budick
相关产品推荐
相关产品推荐

