港股即時行情 WebSocket 動態訂閱實戰:解決 Tick 序號斷層自動補全完整 Python 方案

kalos
·
·
IPFS
·
我一開始開發行情系統時,就是採用「改標的就斷線重連」的寫法,本機測試看似正常,但上線後各種指標失真問題層出不窮。每次重連都會清空本地儲存的序號緩存,搭配網路瞬斷、消費執行緒塞車,即時 Tick 的連續序號會直接中斷;若沒有自動補區間 Tick 的機制,分時均價、累計成交額、K 線、盤口深度全部失準,每次行情劇烈波動都需要維運手動拉歷史資料修復,人力成本極高。

前言

身為金融量化與行情後端開發者,想必都遇過同樣的線上穩定災難:使用者頻繁增刪自選港股、量化策略輪換標的時,若每次異動都直接重連 WebSocket,會大量觸發重連風暴,連帶造成 Tick 資料序號跳脫、數據斷層。

我一開始開發行情系統時,就是採用「改標的就斷線重連」的寫法,本機測試看似正常,但上線後各種指標失真問題層出不窮。每次重連都會清空本地儲存的序號緩存,搭配網路瞬斷、消費執行緒塞車,即時 Tick 的連續序號會直接中斷;若沒有自動補區間 Tick 的機制,分時均價、累計成交額、K 線、盤口深度全部失準,每次行情劇烈波動都需要維運手動拉歷史資料修復,人力成本極高。

後續我分別測試 REST 輪詢、全量重新訂閱兩種傳統方案,最後重構出「單長連線動態增刪訂閱」架構,搭配序號連續性校驗、訊息緩衝、缺失區間自動補全整套機制,從根源解決序號跳脫帶來的資料錯亂。下文會完整分享實作邏輯、可直接執行 Python 程式,還有線上踩坑對策,給同做金融行情開發的夥伴參考。

一、系統核心需求

  1. 單條長存 WebSocket 連線承載數十檔港股,增刪觀察標的時無需中斷連線,維持即時行情持續推送不中斷;

  2. 每筆即時 Tick 攜帶遞增序號,消費端即時校驗連續性,偵測到序號缺口自動呼叫歷史介面拉取缺失區間資料;

  3. 本地維護訂閱標的集合,自動去重、攔截空清單無效指令,避免多餘 Tick 浪費伺服器與本機運算資源;

  4. 內建心跳存活偵測、連線異常自動重連、訊息緩衝亂序兜底,確保各類行情計算指標長期穩定輸出。

二、線上常見資料異常場景

1. 每次重連重置序號,造成大範圍行情斷層

只要增刪標的就銷毀重建 Socket,本地紀錄的上一筆序號會直接清空,新連線從伺服器當下序號重新計數,新舊兩段資料流無法銜接,大量 Tick 永久遺失,回測數據完全失效。

2. 網路抖動、消費速度跟不上,產生小幅序號跳號

營運商封包瞬間遺失、本地處理邏輯壅塞,都會出現序號跳動(例如 1003 直接跳到 1008),中間遺失的 Tick 會直接扭曲成交額、均價、五檔盤口的運算結果。

3. 快速增刪標的引發指令競態,本地與伺服訂閱狀態不一致

使用者連續點擊新增 / 移除自選,多筆訂閱指令並發送出,本地標的集合和伺服器訂閱清單錯位,衍生兩種髒數據:同一 Tick 重複推送、特定標的完全收不到行情。

4. 未設計緩衝層,補全資料與即時 Tick 混雜覆蓋

偵測到序號缺口後只等待歷史資料,新的即時 Tick 持續湧入,新舊資料毫無順序堆疊覆蓋,修復完的分時、K 線圖會出現突兀鋸齒與跳動。

三、核心概念:動態增刪訂閱

動態增刪訂閱指依託單條持續在線的 WebSocket 長連線,透過專屬控制指令攜帶新增、移除標的清單,即時更新伺服器訂閱列表,全程不關閉、不重建 Socket,保障即時行情持續輸送。

此架構和 REST 定時輪詢、異動標的即重連的全量訂閱有本質差異,能大幅降低重連風暴,減少連線資源消耗。

四、業務場景對照驗證表

五、完整可執行 Python 實作程式

import websocket
import json
import threading
import time

# 港股專屬WebSocket行情接入位址
WS_STOCK_URL = "wss://quote.alltick.co/quote-stock-b-ws-api?token=YOUR_TOKEN"

# 本地全域狀態管理
subscriptions = set()  # 儲存當前有效訂閱標的,自動去重避免幽靈訂閱
last_seq = None        # 記錄上一筆Tick序號,用於即時連續性檢核
msg_buffer = []        # 序號缺口緩衝區,補完資料後統一回放消費
ws_app = None

def send_subscribe_action(action: str, code_list: list):
    """發送動態訂閱指令,add新增 / del移除,固定指令ID 22004"""
    global ws_app
    if not ws_app or not ws_app.sock or not ws_app.sock.connected:
        return
    # 攔截空清單無效指令
    if len(code_list) == 0:
        return
    # 標的代碼自動去重
    unique_codes = list(set(code_list))
    sub_frame = {
        "cmd_id": 22004,
        "action": action,
        "code": unique_codes
    }
    ws_app.send(json.dumps(sub_frame))
    # 同步更新本地訂閱集合
    if action == "add":
        for c in unique_codes:
            subscriptions.add(c)
    elif action == "del":
        for c in unique_codes:
            if c in subscriptions:
                subscriptions.remove(c)

def request_missing_tick(start_seq: int, end_seq: int):
    """偵測序號缺口時,呼叫歷史HTTP介面拉取缺失區間Tick"""
    print(f"偵測序號缺口,拉取缺失區間 seq:{start_seq+1} ~ {end_seq-1}")
    # 此處可擴充歷史Tick查詢HTTP邏輯,取回後有序插入緩衝區前方

def check_seq_continuity(current_seq: int) -> bool:
    """即時檢核序號連續性,斷層自動觸發補全機制"""
    global last_seq
    if last_seq is None:
        last_seq = current_seq
        return True
    if current_seq != last_seq + 1:
        request_missing_tick(last_seq, current_seq)
        return False
    last_seq = current_seq
    return True

def on_open(ws):
    """連線建立完成,執行初始批量訂閱"""
    init_codes = ["00700.HK", "9988.HK", "09992.HK"]
    send_subscribe_action("add", init_codes)
    print("WebSocket連線建立完畢,完成港股初始標的訂閱,開始接收即時行情")

def on_message(ws, message):
    """訊息回呼:序號檢核、緩衝儲存、無效報文過濾"""
    global last_seq, msg_buffer
    # 過濾空訊息
    if not message or len(message.strip()) == 0:
        return
    try:
        msg = json.loads(message)
    except Exception:
        return
    # 過濾不含序號、標的代碼的非Tick推送
    if "seq" not in msg or "code" not in msg:
        return
    tick_seq = msg["seq"]
    tick_code = msg["code"]
    # 攔截價格、開盤價皆空的無效行情報文
    if msg.get("price") in (0, None) and msg.get("open") in (0, None):
        return
    # 執行序號連續性檢核
    check_seq_continuity(tick_seq)
    # 有效即時報文存入緩衝,缺口補齊後統一處理
    msg_buffer.append(msg)
    last_seq = tick_seq

def on_error(ws, error):
    print(f"WebSocket連線異常:{error}")

def on_close(ws, close_code, close_msg):
    global last_seq, msg_buffer, subscriptions
    print(f"連線中斷,中斷代碼:{close_code},備註資訊:{close_msg}")
    # 斷線重置所有本地狀態,重連後重新訂閱
    last_seq = None
    msg_buffer.clear()
    subscriptions.clear()

def ws_runner():
    global ws_app
    ws_app = websocket.WebSocketApp(
        WS_STOCK_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, ping_timeout=5)

if __name__ == "__main__":
    # 非同步啟動WebSocket長連線執行緒
    ws_thread = threading.Thread(target=ws_runner, daemon=True)
    ws_thread.start()
    time.sleep(2)
    # 模擬使用者新增自選標的
    send_subscribe_action("add", ["01299.HK"])
    time.sleep(10)
    # 模擬使用者移除訂閱標的
    send_subscribe_action("del", ["01299.HK"])
    while True:
        time.sleep(1)

六、線上故障排查紀錄(現象|偵測方式|解決對策)

1. 高頻 Tick 湧入,訊息緩衝持續堆積溢位

現象:開盤、劇烈波動時段 msg_buffer 長度不斷上漲,消費執行緒 CPU 長期滿載;

偵測:監控緩衝長度日誌、統計單筆 Tick 處理耗時;

對策:獨立非同步執行緒批量消費緩衝,設定緩衝容量上限;超過門檻時拉取行情快照對齊狀態,丟棄過期無效 Tick。

2. 網路假活無報錯,但序號持續斷層

現象:無連線錯誤提示,卻連續多筆 Tick 序號不連續,行情持續缺失;

偵測:監控 ping/pong 心跳回應、統計連續缺口次數;

對策:新增業務層序號逾時偵測,連續 5 次缺口則手動斷線重建,重連後拉取最新快照統一校正資料狀態。

3. 快速增刪標的引發指令競態,本地與伺服訂閱不一致

現象:本地集合已移除的標的,仍持續收到對應 Tick 推送;

偵測:印出 add/del 指令傳送日誌,比對推送標的代碼;

對策:訂閱指令加上執行緒鎖,新增、移除操作循序執行;每筆指令送出後延遲 200ms 驗證推送內容,狀態異常自動修正本地集合。

4. 標的代碼缺少.HK 後綴,訂閱靜默失敗無回報錯

現象:送出訂閱指令後長期收不到對應 Tick,介面不回傳錯誤資訊;

偵測:比對送出 code 欄與介面規範;

對策:本地建立港股代碼白名單,送出新增指令前驗證格式,非法代碼直接攔截不傳送。

七、方案適用邊界

本套動態訂閱架構支援:單條長存 WebSocket 連線內,透過固定指令 ID 動態新增、移除港股標的,全程無需銷毀重建連線,穩定承接即時 Tick;

不支援:多條 WebSocket 之間同步訂閱狀態、一次性完整回溯全量歷史 Tick、規範以外私有互動指令。

八、總結

本篇實作架構基於 AllTick API 的 WebSocket 動態訂閱能力打造,以序號即時檢核搭配區間自動補全,完整解決港股即時 Tick 序號跳脫、資料斷層問題。文中提供可直接複製執行的 Python 程式,同時整理四類線上高頻故障與標準處理流程,適用於量化交易、行情視覺化後端落地。開發時嚴格遵守介面使用邊界,就能有效抑制重連風暴、消除行情指標失真,大幅降低日常維運修復成本。

CC BY-NC-ND 4.0 授权
已推荐到频道:时事・趋势

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