港股即時行情 WebSocket 動態訂閱實戰:解決 Tick 序號斷層自動補全完整 Python 方案
前言
身為金融量化與行情後端開發者,想必都遇過同樣的線上穩定災難:使用者頻繁增刪自選港股、量化策略輪換標的時,若每次異動都直接重連 WebSocket,會大量觸發重連風暴,連帶造成 Tick 資料序號跳脫、數據斷層。
我一開始開發行情系統時,就是採用「改標的就斷線重連」的寫法,本機測試看似正常,但上線後各種指標失真問題層出不窮。每次重連都會清空本地儲存的序號緩存,搭配網路瞬斷、消費執行緒塞車,即時 Tick 的連續序號會直接中斷;若沒有自動補區間 Tick 的機制,分時均價、累計成交額、K 線、盤口深度全部失準,每次行情劇烈波動都需要維運手動拉歷史資料修復,人力成本極高。
後續我分別測試 REST 輪詢、全量重新訂閱兩種傳統方案,最後重構出「單長連線動態增刪訂閱」架構,搭配序號連續性校驗、訊息緩衝、缺失區間自動補全整套機制,從根源解決序號跳脫帶來的資料錯亂。下文會完整分享實作邏輯、可直接執行 Python 程式,還有線上踩坑對策,給同做金融行情開發的夥伴參考。
一、系統核心需求
單條長存 WebSocket 連線承載數十檔港股,增刪觀察標的時無需中斷連線,維持即時行情持續推送不中斷;
每筆即時 Tick 攜帶遞增序號,消費端即時校驗連續性,偵測到序號缺口自動呼叫歷史介面拉取缺失區間資料;
本地維護訂閱標的集合,自動去重、攔截空清單無效指令,避免多餘 Tick 浪費伺服器與本機運算資源;
內建心跳存活偵測、連線異常自動重連、訊息緩衝亂序兜底,確保各類行情計算指標長期穩定輸出。
二、線上常見資料異常場景
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 程式,同時整理四類線上高頻故障與標準處理流程,適用於量化交易、行情視覺化後端落地。開發時嚴格遵守介面使用邊界,就能有效抑制重連風暴、消除行情指標失真,大幅降低日常維運修復成本。
喜欢我的作品吗?别忘了给予支持与赞赏,让我知道在创作的路上有你陪伴,一起延续这份热忱!