WebSocket数据流断收后如何重启asyncio协程
问题核心原因
你遇到的无征兆断流基本都是WebSocket半开连接导致的:网络中间节点(NAT网关、防火墙、运营商链路)超时回收了空闲连接,但TCP层没有收到连接断开的控制报文,从程序视角看连接还是“正常”的,await 调用既不会抛出异常,也收不到新数据,部分场景下ccxtpro会直接返回None或空响应。
原代码存在三个明显缺陷,无法应对这种场景:
- 仅捕获了主动抛出的异常,没有对空响应、长时间无数据的静默断流做检测
- 任务异常退出后直接关闭客户端,没有重连逻辑
- 没有给流调用加超时兜底,半开连接会永久卡住协程
修复方案
核心改动逻辑
- 给所有
watch_*调用加超时检测,超过阈值没收到数据直接判定断流 - 增加空响应校验,只要返回
None/空结构就触发重连 - 每个数据流独立维护客户端和重连循环,单个流出错不影响其他流
- 重连加指数退避逻辑,避免高频请求被交易所限流
- 所有资源销毁逻辑放在
finally块,避免连接、协程泄漏
修改后可直接运行的代码
# tasks.py 保持原有逻辑无需改动 @app.task(bind=True, name='Start websocket loops') def start_ws_loops(self): ws_loops()
# methods.py import asyncio import logging from asyncio import gather, wait_for, TimeoutError log = logging.getLogger(__name__) # 配置项可根据业务调整 STREAM_TIMEOUT = 360 # 单流最大空闲时间,单位秒,5分钟K线设360s(6分钟)留冗余 RECONNECT_BASE_DELAY = 1 # 首次重连等待时间 RECONNECT_MAX_DELAY = 30 # 重连最大等待时间 def ws_loops(): async def method_loop(client, method, private, args): while True: # 给每次监听调用加超时,卡住超过阈值直接抛错 if private: response = await wait_for( getattr(client, method)(), timeout=STREAM_TIMEOUT ) else: response = await wait_for( getattr(client, method)(**args), timeout=STREAM_TIMEOUT ) # 空值直接判定为断流 if response is None or (isinstance(response, (list, dict)) and len(response) == 0): raise ConnectionError(f"Stream {method} {args} return empty response") # 正常业务处理 if method in ('watchMyTrades', 'watchOrders', 'watch_ohlcv'): do_stuff(response) async def single_stream_worker(stream_config): """单流独立工作协程,自带无限重连,和其他流完全隔离""" exid = stream_config['exid'] wallet = stream_config['wallet'] method = stream_config['method'] private = stream_config['private'] args = stream_config['args'] reconnect_delay = RECONNECT_BASE_DELAY exchange = Exchange.objects.get(exid=exid) while True: client = None try: parameters = {'enableRateLimit': True, 'newUpdates': True} if private: log.info(f'Init private stream: {method} {args}') client = exchange.get_ccxt_client_pro(parameters, wallet=wallet, account=args['account']) else: log.info(f'Init public stream: {method} {args}') client = exchange.get_ccxt_client_pro(parameters, wallet=wallet) # 启动监听,正常情况下会永久阻塞在这里收数据 await method_loop(client, method, private, args) except Exception as e: log.warning(f"Stream {method} {args} down, retry after {reconnect_delay}s, err: {str(e)}") await asyncio.sleep(reconnect_delay) # 指数退避,直到达到最大等待时间 reconnect_delay = min(reconnect_delay * 2, RECONNECT_MAX_DELAY) finally: # 强制关闭旧客户端释放资源 if client is not None: try: await client.close() except Exception: pass # 连接成功后重置重连等待时间 else: reconnect_delay = RECONNECT_BASE_DELAY async def main(): stream_list = [] private_methods = ['watchMyTrades', 'watchOrders'] public_methods = ['watch_ohlcv'] for exid in ['binance']: for wallet in ['spot', 'future']: # 组装私有流配置 for method in private_methods: for account in ['user1', 'user2']: stream_list.append(dict( exid=exid, wallet=wallet, method=method, private=True, args=dict(account=account) )) # 组装公共流配置 for method in public_methods: for symbol in ['ETH/USD', 'BTC/USD', 'BNB/USD']: stream_list.append(dict( exid=exid, wallet=wallet, method=method, private=False, args=dict(symbol=symbol, timeframe='5m', limit=1) )) # 所有流独立启动,单个流出错不影响其他流 tasks = [asyncio.create_task(single_stream_worker(config)) for config in stream_list] await gather(*tasks) # 用官方推荐的asyncio.run启动,自动管理事件循环生命周期,不需要手动传loop参数 asyncio.run(main())
额外注意事项
- ccxtpro自带的重连逻辑仅能处理明确抛出的连接错误(比如交易所主动发断开帧、DNS解析失败、鉴权过期),半开连接属于传输层静默故障,上层库无法感知,必须自己加超时检测。
- 超时阈值要和订阅的流匹配:5分钟周期K线设300360秒超时即可,成交、订单簿这类实时推送的流超时设3060秒即可,阈值太短容易误判正常空闲,太长会导致断流后恢复慢。
- 如果用Celery启动该任务,worker必须使用
solo池或者gevent协程池,默认的prefork多进程池和asyncio事件循环存在兼容性问题,会随机出现协程卡住、不执行的问题。
内容的提问来源于stack exchange,提问作者Florent
相关产品推荐
相关产品推荐

