规避回测偏差:个股停牌 Tick 快照完整采集方案

用户头像sh_****447dvu
2026-07-28 发布

一、研究背景与问题现象

在基于实时 Tick 数据搭建量化回测、策略仿真系统时,个股停牌复牌场景极易引发数据时序缺陷,直接干扰模型收益、波动率、开平仓信号等核心指标计算。

常规直连行情 WebSocket 的简易实现存在明显缺陷:当标的长期停牌后恢复交易,K 线时序出现空白断层,回测数据集缺失停牌区间交易状态标记,导致历史拟合、样本外验证结果失真,策略可靠性无法客观评估。

早期调试阶段采用粗暴兜底逻辑:长时间无 Tick 推送即判定链路失效,反复重建 WebSocket 连接、批量拉取历史 K 线补全数据。该方案会带来三重负面影响:接口请求频次激增、带宽资源冗余消耗、数据库产生大量重复快照记录,数据清洗环节额外增加大量校验成本。

经多轮实盘数据回放、离线回测验证,形成一套标准化时序修复流程:依托标准化 WebSocket 订阅指令,解析报文内置status交易状态字段,单长连接动态管理标的订阅,无需频繁销毁重建链路,不虚构停牌区间无成交 Tick,自动补全停牌至复牌完整时序快照,稳定支撑量化模型的数据输入要求。

二、核心逻辑定义

停牌复牌快照时序修复:标的复牌首条 Tick 到达后,读取本地持久化的停牌状态缓存,校验停牌周期内历史快照完整性,仅补充交易状态标记字段,串联停牌前静态价格快照与复牌实时数据流。

区别于两类低效实现:

  1. 无需销毁、重建 WebSocket 连接,规避大量短连接引发的重连风暴与时序断点;
  2. 不依赖定时 REST 轮询批量拉取历史 K 线,减少无效数据请求,降低回测数据加载耗时。

三、量化研究高频场景与标准化处理对照表

应用场景 量化研究痛点 订阅配置(cmd_id/action/code) 回测数据复核基准
初始化订阅含停牌标的 无法区分接口断连 / 标的停牌,误判数据缺失触发全量历史补拉,拉长回测初始化耗时 cmd_id=22004,action=subscribe,解析 status 字段 WebSocket 握手完成,同步缓存每只标的交易状态至本地时序库
盘中临时停牌 短时无 Tick 推送频繁重连,回测回放时出现大量重复行情片段 维持单条长连接持续订阅,识别suspend停牌标识 本地记录停牌起止时间,留存停牌前最后一笔成交快照
长期停牌后复牌 复牌 Tick 直接入库,时序空白造成回测 K 线缺口,策略信号计算偏移 复用原有 WebSocket 通道,自动触发快照时序修复流程 匹配停牌起止时间,补全区间状态记录,保证时间轴连续无断点
重复订阅停牌标的 重复下发订阅指令,快照数据冗余,回测样本重复计数 订阅前本地集合去重拦截上行请求 缓存校验已订阅标的,阻断重复订阅指令下发
网络断连恰逢标的复牌 重连丢失停牌状态记录,前后行情时序割裂,回测分段数据无法拼接 重连自动恢复订阅,批量拉取标的基础交易状态 重连完成后执行一次全量历史快照完整性校验

四、量化回测四大典型数据缺陷及修复方案

1. 将停牌无 Tick 等同于数据丢失,循环请求历史接口

现象:标的停牌期间无成交报文,程序持续调用批量 K 线接口,回测数据集加载效率大幅下降,接口额度快速消耗。

数据检测方式:解析每条 Tick 内置status字段,区分「交易所交易暂停」与「链路异常断开」两类无数据区间。

标准化处理规则:仅当标的状态为正常交易、且长时间无数据推送时执行链路重连;停牌状态下关闭历史数据批量补拉逻辑。

2. 复牌 Tick 入库未关联停牌元数据,时序链条断裂

现象:数据库仅留存停牌前快照、复牌首条 Tick,中间时段无状态记录,回测绘图、指标计算出现固定缺口。

数据检测方式:提取标的完整时间序列,比对停牌结束时间与复牌 Tick 时间戳连续性。

标准化处理规则:独立构建交易状态元数据表,存储每只标的停牌起止时间、停牌前收盘价;复牌数据入库时关联元数据表,写入时序标记填补空白区间。

3. 多标的同步复牌,并行修复引发数据库写入竞态

现象:多标的同日复牌,多线程并行执行修复逻辑,同一标的 + 停牌周期生成多条重复状态记录,回测样本重复统计。

数据检测方式:按code+suspend_start分组统计记录条数,重复条目判定为竞态问题。

标准化处理规则:单标的快照修复逻辑串行执行;数据库设置code+suspend_start联合唯一索引,杜绝重复存储。

4. 跨品类混用 WebSocket 地址,停牌状态字段无法解析

现象:股票标的使用加密品类 WebSocket 地址订阅,报文缺失status字段,无法识别停牌、复牌,回测持续出现时序漏洞。

数据检测方式:核对接入域名,股票品类需使用独立专用 WebSocket 链路。

标准化处理规则:代码层做品类路由隔离,股票行情请求强制路由至股票专属 WSS 地址,拦截跨品类错误请求。

五、方案适用边界说明

本流程基于标准订阅指令cmd_id=22004构建,支持单条活跃 WebSocket 连接内动态增删标的、修复停牌时序,适配离线回测、实盘仿真、策略样本采集场景;存在两处固定边界,研究开发阶段需提前适配:

  1. 无法在多条独立 WebSocket 连接之间同步个股停牌状态,多进程回测需统一中心化状态缓存;
  2. 不会自动生成停牌期间虚拟 Tick 成交数据,仅补充交易状态标记,不支持人为构造模拟行情用于回测。

六、Python 完整可运行代码(适配量化数据采集)

import websockets
import asyncio
import json
from datetime import datetime

# 股票行情专用WSS链路,参考官方接口文档
WSS_STOCK_URL = "wss://quote.alltick.co/quote-stock-b-ws-api?token=YOUR_TOKEN"

class StockQuoteDataCollector:
    def __init__(self):
        self.ws = None
        self.subscriptions = set()
        # 本地时序缓存:key=股票code,存储停牌时间、停牌前价格、交易状态
        self.stock_status_cache = {}

    async def send_subscribe(self, action: str, code_list: list):
        if not code_list:
            return
        payload = {
            "cmd_id": 22004,
            "action": action,
            "code": code_list
        }
        await self.ws.send(json.dumps(payload))
        if action == "subscribe":
            [self.subscriptions.add(c) for c in code_list]
        elif action == "unsubscribe":
            [self.subscriptions.discard(c) for c in code_list]

    def check_resume_repair(self, code: str, curr_status: str, trade_time: str):
        """量化数据核心:检测标的复牌,启动时序快照修复逻辑"""
        cache_info = self.stock_status_cache.get(code)
        if not cache_info:
            return
        old_status = cache_info["status"]
        # 交易状态由停牌切换为正常,判定复牌触发修复
        if old_status == "suspend" and curr_status == "normal":
            print(f"标的{code}复牌,执行回测时序快照校验")
            self.repair_snapshot_timeline(code, cache_info["suspend_start"], trade_time)
            self.stock_status_cache[code]["status"] = "normal"

    def repair_snapshot_timeline(self, code, suspend_start, resume_time):
        """持久化停牌区间时序记录,保障回测时间轴完整"""
        repair_record = {
            "code": code,
            "suspend_start": suspend_start,
            "resume_time": resume_time,
            "pre_suspend_price": self.stock_status_cache[code]["last_price"],
            "status": "suspend_repaired"
        }
        # save_backtest_snapshot(repair_record) 替换为量化系统持久化逻辑
        print("写入停牌区间时序修复记录,用于离线回测", repair_record)

    async def on_open(self):
        # 初始化回测观测标的池
        init_codes = ["NASDAQ:AAPL", "HKEX:00700"]
        await self.send_subscribe("subscribe", init_codes)
        print("股票WebSocket链路建立,完成回测标的初始订阅")

    async def on_message(self, raw_msg):
        if not raw_msg:
            return
        try:
            data = json.loads(raw_msg)
            tick_data = data.get("data", {})
            code = tick_data.get("code")
            price = tick_data.get("price")
            trade_time = tick_data.get("time")
            status = tick_data.get("status", "normal")

            # 空值过滤,剔除无效数据,避免污染回测数据集
            if not code or price in (None, 0) or not trade_time:
                return

            # 更新本地标的状态缓存
            if code not in self.stock_status_cache:
                self.stock_status_cache[code] = {}
            self.stock_status_cache[code]["last_price"] = price
            self.stock_status_cache[code]["status"] = status
            if status == "suspend" and "suspend_start" not in self.stock_status_cache[code]:
                self.stock_status_cache[code]["suspend_start"] = trade_time

            # 校验是否触发复牌时序修复
            self.check_resume_repair(code, status, trade_time)
            print(f"Tick采集 | {code} 成交价:{price} 交易状态:{status}")
        except Exception as e:
            print("Tick报文解析异常,丢弃本条脏数据", str(e))

    async def on_error(self, err):
        print("WebSocket链路异常,中断数据采集:", err)

    async def on_close(self):
        print("股票WebSocket链路关闭,暂停行情采集")

    async def connect(self):
        try:
            async with websockets.connect(
                WSS_STOCK_URL,
                ping_interval=10
            ) as ws:
                self.ws = ws
                await self.on_open()
                while True:
                    msg = await ws.recv()
                    await self.on_message(msg)
        except Exception as e:
            await self.on_error(e)
            await self.on_close()

async def run_backtest_collector():
    collector = StockQuoteDataCollector()
    task = asyncio.create_task(collector.connect())
    await task

if __name__ == "__main__":
    asyncio.run(run_backtest_collector())

七、研究总结

量化策略的回测有效性高度依赖时序连续、状态完整的行情数据集,停牌复牌带来的时序断层属于极易被忽略、但会系统性扭曲模型输出的底层数据缺陷。

完整的数据修复体系需要三层逻辑协同:实时 Tick 数据流状态监听、本地标的交易状态缓存、历史快照时序完整性校验,三层配合可从源头消除 K 线缺口、指标计算偏移、样本重复统计等影响回测可信度的问题。

若需要搭建覆盖股票、外汇、贵金属、加密货币全品类的标准化行情采集工具,快速落地本套时序修复逻辑,可采用 AllTick API。统一规范的 WebSocket 订阅报文、多语言完整示例代码,能够减少多品类特殊交易场景的数据适配、调试验证工作量,稳定支撑离线回测、实盘仿真、因子挖掘等量化研究工作。

评论