多股监控卡顿?动态订阅稳住 A 股实时 API

场景描述

我长期使用 A 股实时行情 API 搭建个人量化高频交易程序,日常同时监控数十只沪深个股,在线运行时经常遭遇两类严重影响策略回测、实盘一致性的线上问题。第一种是频繁切换股票标的时,直接关闭、重建 WebSocket 长连接,短时间批量发起握手请求,触发接口限流,行情推送卡顿,策略反复重复计算持仓、交易信号;第二种是弱网环境下 Socket 假活,tick 数据流长时间断档才触发断开回调,本地订阅列表与服务端不同步,停牌个股行情漏收,最终导致回测结果和实盘交易收益出现明显偏差。之前尝试过 REST 轮询、多连接并行订阅两种方案,要么行情延迟过高,要么连接数量超限被接口限制。接入 AllTick A 股实时行情 API 后,采用单连接动态增减订阅的 WebSocket 架构,彻底降低了连接波动带来的量化策略误差。

开发与量化核心需求

作为量化开发者,基于 A 股实时行情 API 搭建行情订阅系统,有三项不可妥协的开发诉求:

  1. 切换监控个股无需销毁重建长连接,杜绝批量重连引发的握手限流、重连风暴;
  2. 本地订阅状态与行情服务端实时对齐,消除幽灵订阅、停牌复牌 tick 漏采集问题;
  3. 原生支持心跳保活逻辑,网络抖动可快速识别回调,行情数据统一区分正常交易 / 停牌状态,适配 A 股回测数据清洗流程。

A 股行情 API 高频开发痛点

  1. 多连接订阅方案:每新增一只 A 股个股新建一条 WebSocket,标的池扩容后连接数快速触碰接口阈值,断连后批量重连抢占带宽,tick 回调堆积阻塞量化策略计算;
  2. 销毁重连式订阅:切换个股直接 close 连接再重新握手,每次初始化需要全量推送股票 code 列表,停牌标的重复订阅,无效空成交量数据占用程序内存;
  3. 缺少本地订阅状态校验逻辑:接收 tick 数据不校验标的是否在本地订阅池,接口返回停牌零成交数据时,量化策略误触发买卖信号;
  4. 无增量订阅指令:仅支持一次性全量下发股票清单,无法单独新增、取消单只个股,盘中调仓切换标的开发成本高。

落地方案:AllTick A 股实时行情 API cmd_id=22004 单连接动态订阅模型

概念定义

动态增减订阅指在同一条存活 WebSocket 长连接内,通过下发 cmd_id=22004 指令,搭配 action 字段区分新增 / 取消股票 code 列表完成标的切换;和销毁重建连接、REST 轮询拉取行情两种方式完全区分,全程复用握手、心跳链路,不产生额外握手开销,从根源规避重连风暴。

可复核落地对照表

表格

应用场景高频开发痛点AllTick A 股 API 动态参数配置(cmd_id/action/code)复核校验基准
程序启动初始化订阅首次连接批量加载沪深个股,一次性注册全部监控标的 codecmd_id=22004,action="add",code=[A 股标的列表]连接 on_open 回调中执行,本地 subscriptions 集合同步写入全部股票 code
盘中新增监控个股盘中临时加入个股观察池,关闭重连会中断当前 tick 数据流cmd_id=22004,action="add",code=[单只 / 多只 A 股 code]发送指令前校验 code 不在本地集合,去重后下发,收到 tick 再写入缓存集合
盘中取消监控个股个股停牌、完成清仓,无需持续接收行情数据cmd_id=22004,action="del",code=[待取消 A 股 code]指令下发成功后,本地 subscriptions 同步移除对应标的,过滤后续无关 tick
边界:重复订阅同一 A 股标的业务逻辑多次触发新增订阅,服务端重复推送 tick 造成回调冗余cmd_id=22004,action="add",code=[已存在股票 code]本地集合先判断标的已存在则直接跳过,不向行情 API 重复发送订阅帧
边界:空列表订阅指令代码异常传入空标的数组,触发 API 报错cmd_id=22004,action="add/del",code=[]本地前置校验列表长度,空数组直接拦截,不发起 WebSocket 指令

Python 完整可运行代码(适配 AllTick A 股实时行情 API,动态订阅逻辑锁)

python

运行

import websocket
import json
import time

# 端点/订阅格式:AllTick 官方 API Docs(WebSocket 地址说明)
# A股股票专用WSS行情地址
STOCK_WSS_URL = "wss://quote.alltick.co/quote-stock-b-ws-api?token=YOUR_TOKEN"
# 外汇/加密/商品通用WSS行情地址
COMMON_WSS_URL = "wss://quote.alltick.co/quote-b-ws-api?token=YOUR_TOKEN"

# 本地订阅状态存储,实现标的自动去重
subscriptions = set()

def send_subscribe_frame(ws, action, code_list):
    """统一封装22004订阅指令,单连接内动态增删A股标的,全程不重建连接"""
    if not code_list or len(code_list) == 0:
        # 拦截空列表指令,避免行情API抛出异常
        return
    # 本地前置去重校验
    target_codes = []
    for code in code_list:
        if action == "add" and code not in subscriptions:
            target_codes.append(code)
        elif action == "del" and code in subscriptions:
            target_codes.append(code)
    if len(target_codes) == 0:
        return
    # 组装订阅指令帧,严格遵循AllTick A股实时行情API文档cmd_id规范
    frame = {
        "cmd_id": 22004,
        "action": action,
        "code": target_codes
    }
    ws.send(json.dumps(frame))
    # 同步更新本地订阅缓存集合
    if action == "add":
        subscriptions.update(target_codes)
    elif action == "del":
        for c in target_codes:
            subscriptions.discard(c)

def on_open(ws):
    """WebSocket连接建立,批量初始化订阅沪深A股标的"""
    init_codes = ["600036.SH", "000001.SZ"]
    send_subscribe_frame(ws, "add", init_codes)
    print("行情长连接建立,完成A股初始标的订阅")

def on_message(ws, message):
    """行情回调,过滤无效数据、区分A股停牌状态,适配量化回测逻辑"""
    if not message:
        return
    data = json.loads(message)
    tick_code = data.get("code", "")
    # 过滤不在订阅池内的幽灵行情数据
    if tick_code not in subscriptions:
        return
    # A股停牌数据判定逻辑:当日成交量为0、无有效成交价格
    vol = data.get("volume", 0)
    close_price = data.get("close", 0)
    if vol == 0 and close_price == 0:
        # 停牌个股仅记录状态,不送入指标计算、不生成交易信号
        print(f"A股标的{tick_code}当前停牌,跳过量化指标计算")
        return
    # 正常交易tick,送入量化策略计算模块
    print(f"A股正常行情:{tick_code},现价{close_price},成交量{vol}")

def on_error(ws, error):
    print(f"A股行情WebSocket连接异常:{str(error)}")

def on_close(ws, close_code, close_msg):
    print(f"A股行情长连接断开,状态码:{close_code},信息:{close_msg}")
    # 清空本地订阅缓存,重连后重新初始化订阅列表
    subscriptions.clear()

if __name__ == "__main__":
    ws_app = websocket.WebSocketApp(
        STOCK_WSS_URL,
        on_open=on_open,
        on_message=on_message,
        on_error=on_error,
        on_close=on_close
    )
    # 10秒自动心跳ping,快速识别Socket假活异常
    ws_app.run_forever(ping_interval=10)

开发避坑记录

  1. 现象:A 股高频 Tick 海量推送,本地回调队列堆积、程序内存持续上涨检测:打印单条 tick 回调处理耗时,耗时持续攀升、内存占用无回落兜底:on_message 回调轻量化过滤,停牌、非订阅标的数据直接丢弃,不进入量化计算队列;使用异步任务池隔离指标计算逻辑,防止阻塞 WebSocket 主线程。
  2. 现象:网络波动导致 Socket 假活,心跳周期内无 on_close 回调,A 股行情静默断流检测:连续 10 个心跳周期未收到 API 推送 tick 数据兜底:本地新增行情接收计时器,超时主动关闭 Socket,外层重连逻辑重建链路,清空订阅集合后重新下发 A 股标的列表。
  3. 现象:快速增删 A 股订阅产生竞态,本地订阅集合与服务端指令不同步,出现幽灵订阅检测:日志打印订阅缓存,已执行 del 取消的标的仍持续接收 tick 行情兜底:所有 add/del 指令前置本地集合校验,指令下发完成同步更新缓存;接收 tick 时二次校验标的 code,过滤残留无关行情。
  4. 现象:A 股标的 code 格式不规范(缺失.SH/.SZ 交易所后缀),订阅静默失败,无报错但无行情推送检测:下发订阅指令后长时间无对应个股 tick 返回兜底:统一维护 A 股标的 code 标准化映射表,订阅前统一补齐交易所后缀,对齐 AllTick A 股产品代码文档规范。

功能边界声明

这套动态订阅模型支持:单条活跃 WebSocket 长连接内,通过 cmd_id=22004 指令动态新增、取消任意数量 A 股标的 code;不支持:跨多条 WebSocket 连接同步订阅状态、调用 API 回溯历史 tick 数据、使用非 cmd_id=22004 的私有订阅控制指令。

总结

针对 A 股高频量化开发中普遍存在的 WS 断连、重复计算、订阅不同步、停牌数据干扰等痛点,单连接动态订阅是一套低成本、易落地的优化方案。依托单长连接复用 + 增量订阅指令,既能规避批量重连引发的限流与计算冗余,又能通过本地状态校验统一处理停牌行情,缩小回测与实盘的结果偏差。整套逻辑可直接复用文中 Python 代码快速落地,借助 AllTick API 标准的 WebSocket 订阅能力,无需额外封装复杂中间层,就能稳定承载多标的实时行情监控需求,大幅降低量化开发者的行情链路维护成本。

参考文档:https://apis.alltick.co/
GitHub:https://github.com/alltick/alltick-realtime-forex-crypto-stock-tick-finance-websocket-api

0
0
0
0
评论
未登录
暂无评论