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

WebSocket数据流断收后如何重启asyncio协程

问题核心原因

你遇到的无征兆断流基本都是WebSocket半开连接导致的:网络中间节点(NAT网关、防火墙、运营商链路)超时回收了空闲连接,但TCP层没有收到连接断开的控制报文,从程序视角看连接还是“正常”的,await 调用既不会抛出异常,也收不到新数据,部分场景下ccxtpro会直接返回None或空响应。

原代码存在三个明显缺陷,无法应对这种场景:

  • 仅捕获了主动抛出的异常,没有对空响应、长时间无数据的静默断流做检测
  • 任务异常退出后直接关闭客户端,没有重连逻辑
  • 没有给流调用加超时兜底,半开连接会永久卡住协程

修复方案

核心改动逻辑

  1. 给所有watch_*调用加超时检测,超过阈值没收到数据直接判定断流
  2. 增加空响应校验,只要返回None/空结构就触发重连
  3. 每个数据流独立维护客户端和重连循环,单个流出错不影响其他流
  4. 重连加指数退避逻辑,避免高频请求被交易所限流
  5. 所有资源销毁逻辑放在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 06:57:22