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

如何立即终止所有异步运行任务?附Python WebSocket代码问题

问题分析

你的代码核心问题有两个:

  1. await websocket.recv()是无阻塞上限的调用,Bybit的WebSocket会先返回订阅确认消息(而非目标数据),后续还可能发送心跳包,这些都会被你的代码接收但不会计入ws_responses,导致循环一直等待,直到真正的深度数据到来(可能延迟40秒)。
  2. 没有超时机制,一旦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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 18:05:18