WebSocket 行情實戰:修復復牌時序斷層,維持量化回測數據連續性

kalos
·
·
IPFS
在搭建基於 Tick 即時數據的回測框架、策略模擬平台時,個股停牌與復牌所造成的時間軸斷裂,往往會直接干擾收益曲線、波動率、開平倉訊號等核心因子的計算準確度。若僅使用最基礎的 WebSocket 行情直連邏輯,長期停牌標的恢復交易後,K 線時間軸會出現空白區間;資料庫缺少停牌期間的狀態標記,會造成歷史擬合、樣本外測試的結果失真,無法客觀判斷策略的穩定性。

前言:量化研究中容易被忽視的底層數據缺陷

在搭建基於 Tick 即時數據的回測框架、策略模擬平台時,個股停牌與復牌所造成的時間軸斷裂,往往會直接干擾收益曲線、波動率、開平倉訊號等核心因子的計算準確度。

若僅使用最基礎的 WebSocket 行情直連邏輯,長期停牌標的恢復交易後,K 線時間軸會出現空白區間;資料庫缺少停牌期間的狀態標記,會造成歷史擬合、樣本外測試的結果失真,無法客觀判斷策略的穩定性。

一開始我採取簡易兜底機制:只要長時間未收到 Tick 推送,便判定連線中斷,重複重建 WebSocket、批量拉取歷史 K 線填補空缺。這套作法會衍生三項問題:重複請求佔用資源、資料庫堆積大量重複快照、數據清洗階段耗費額外人力與運算成本。

透過多次離線回測與歷史行情回放驗證,我整理一套標準化時序修復流程:依據 WebSocket 標準訂閱指令解析報文中的status交易狀態欄位,以單一長連線動態管理觀測標的,無需頻繁重連,亦不會虛構停牌期間不存在的成交 Tick,自動補齊停牌至復牌的完整快照鏈路,穩定供給量化模型高品質輸入數據。

一、核心機制定義:停牌復牌快照時序修復

當標的復牌後第一筆 Tick 抵達時,程式讀取本地持久化儲存的停牌緩存資訊,校驗停牌區間內所有歷史快照的完整性,僅補充交易狀態標記,串聯停牌前靜態價格快照與復牌即時數據流。

相較兩種常見低效實作,本機制具備明顯優勢:

  1. 無需銷毀、重建 WebSocket 連線,避免大量短連線引發的重連風暴與時間軸斷點;

  2. 不依靠定時 REST 輪詢批量拉取歷史 K 線,減少多餘請求,縮短回測數據載入時間。

二、量化研究常見場景與標準處理對照表

三、四類常見數據瑕疵與標準修復邏輯

1. 將停牌無 Tick 視為資料遺失,持續呼叫歷史介面

現象:標的停牌期間無任何成交報文,程式持續請求批量 K 線,回測載入效率大幅下滑。

辨識方式:解析每筆 Tick 內建status欄位,區分「交易所暫停交易」與「真實連線中斷」兩種無數據區間。

處理規則:僅當標示為正常交易、且長時間無資料時,才執行重連;停牌狀態下關閉歷史數據補拉機制。

2. 復牌 Tick 寫庫時未綁定停牌元數據,時序鏈路斷裂

現象:資料庫僅留存停牌前快照與復牌第一筆 Tick,中間無狀態紀錄,繪圖、指標計算皆出現固定缺口。

辨識方式:提取單一標的完整時間序列,比對停牌結束時間與復牌 Tick 時間戳是否連續。

處理規則:獨立建置交易狀態元數據表,儲存各標的停牌起止時間、停牌前收盤價;復牌資料寫入時綁定對應紀錄,填入時序標記填補空白。

3. 多標的同步復牌,並行修復引發資料庫寫入競爭

現象:多檔股票同日復牌,多執行緒同步執行修復邏輯,同一標的、同一停牌區間產生多筆重複紀錄,樣本重複統計。

辨識方式:以code搭配suspend_start分組統計筆數,重複條目即代表競爭問題。

處理規則:單一標的快照修復採循序執行;資料庫建立code+suspend_start複合唯一索引,阻擋重複儲存。

4. 跨品類混用 WebSocket 位址,無法解析停牌標記

現象:股票使用加密貨幣專用連線訂閱,報文缺少status欄位,無法判斷停牌、復牌,回測持續出現時序漏洞。

辨識方式:檢查連線網域,股票需使用獨立專屬 WSS 通道。

處理規則:程式內建立品類路由隔離,股票行情強制導向專屬位址,攔截跨品類錯誤請求。

四、方案適用邊界

整套流程以標準訂閱指令cmd_id=22004為核心,支援單一 WebSocket 連線內動態增減標的、修復停牌時序,適用離線回測、策略模擬、因子採集等場景;實作前需留意兩項限制:

  1. 無法在多條獨立 WebSocket 連線間同步停牌狀態,多行程式回測需搭配中心化狀態緩存;

  2. 不會自動生成停牌期間虛擬 Tick,僅補充交易狀態標記,不支援人為建構模擬行情。

完整可執行 Python 程式(量化數據採集適用)

import websockets
import asyncio
import json
from datetime import datetime

# 股票專屬WSS通道,參考官方API文件
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=股票代碼,存停牌時間、停牌前價格、交易狀態
        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())

研究總結

量化策略的可信度建立在連續、完整的行情數據集之上,停牌復牌造成的時間軸缺口屬於底層隱性瑕疵,會系統性偏離因子與回測結果。一套完備的數據處理架構,必須同時具備即時流狀態監控、本地標的緩存、歷史快照校驗三層邏輯,才能從根源消除 K 線缺口、指標偏移、樣本重複等問題。

若需要建置橫跨股票、外匯、貴金屬、加密貨幣的多品類行情採集工具,快速落地這套時序修復邏輯,可使用 AllTick API。統一規格的 WebSocket 訂閱報文與多語言範例,能大幅減少各市場特殊交易場景的調試成本,穩定支援因子挖掘、離線回測、實盤模擬等量化研究工作。

CC BY-NC-ND 4.0 授权

喜欢我的作品吗?别忘了给予支持与赞赏,让我知道在创作的路上有你陪伴,一起延续这份热忱!