如果你曾負責過跨國資產配置或大中華區股票量化策略開發,一定碰過這種讓人抓狂的場景:回測時,你在本地硬碟讀取預先清洗好的分K歷史資料,各項指標與回報率漂亮得無可挑剔;然而一把策略部署到模擬盤或實盤,切換到實時串流時,整個程式卻彷彿人格分裂。
訊號觸發延遲、指標算出來的值與回測兜不攏,甚至連代碼寫法都得配合兩套完全不同的接口改得面目全非。問題的癥結往往不在你的數學模型,而在於底層「歷史資料」與「實時行情」被硬生生割裂成了兩座孤島。
企業金融分析師的痛點:資料不合一帶來的效率黑洞
在日常量化交易維運與資料清洗管線中,這種架構缺陷會引發一連串連鎖反應:
- 重複撰寫介面轉接層: 為了讀取歷史分線寫一套 REST Parser,為了監聽即時報價又得寫一套 WebSocket Handler,耗費大量精力在搬磚而非策略優化。
- 分界點重複與污染: 策略啟動當下,REST 補進來的最後一根分K往往尚未走完,而 WebSocket 即時推送的第一筆 Tick 又恰好落在同一分鐘內,稍不留神就會產生資料重複或覆蓋錯誤。
- 交易時段邏輯誤判: A股與台股或美股不同,每天 11:30 至 13:00 有整整一個半小時的午休休市,缺乏狀態管制的連線程式極易把這種正常靜默誤判為斷線重連。
雙軌合一架構:REST補歷史,WebSocket接實時
要徹底根治這個問題,核心想法非常單純:在架構底層將「過去」與「現在」職責拆開,對上層策略卻提供絕對統一的高階呼叫。
歷史資料走 REST 一次性拉取整批,實時報價走 WebSocket 由伺服器主動發送。在實作跨市場資料整合時,選用具備完整規格的標準數據源(例如整合度高的 AllTick API)能讓連線協議極度簡化。
import asyncio
import json
import os
import uuid
from dataclasses import dataclass
import requests
import websockets
TOKEN = os.environ["ALLTICK_API_TOKEN"]
REST_URL = "https://quote.alltick.co/quote-stock-b-api/kline"
WS_URL = "wss://quote.alltick.co/quote-stock-b-ws-api?token=" + TOKEN
@dataclass
class Bar:
ts: int # 秒級時間戳,K線起點
open: float
high: float
low: float
close: float
volume: float
class AStockData:
"""策略只調用 history() 和 stream_ticks(),不關心底層細節"""
def history(self, code, num=500):
query = {
"trace": str(uuid.uuid4()),
"data": {
"code": code,
"kline_type": 1, # 1代表1分鐘K線
"kline_timestamp_end": 0, # 股票只支援0:從最新交易日往前取
"query_kline_num": num, # 單次最多500根
"adjust_type": 0, # 目前僅支援0(不復權)
},
}
resp = requests.get(
REST_URL,
params={"token": TOKEN, "query": json.dumps(query)},
timeout=10,
)
resp.raise_for_status()
body = resp.json()
if body.get("ret") != 200:
raise RuntimeError(body.get("msg"))
bars = [
Bar(
ts=int(k["timestamp"]),
open=float(k["open_price"]),
high=float(k["high_price"]),
low=float(k["low_price"]),
close=float(k["close_price"]),
volume=float(k["volume"]),
)
for k in body["data"]["kline_list"]
]
return sorted(bars, key=lambda b: b.ts)
async def stream_ticks(self, codes):
subscribe = {
"cmd_id": 22004,
"seq_id": 1,
"trace": str(uuid.uuid4()),
"data": {"symbol_list": [{"code": c} for c in codes]},
}
heartbeat = {"cmd_id": 22000, "seq_id": 1, "trace": "heartbeat", "data": {}}
while True: # 斷線後自動重連
try:
async with websockets.connect(WS_URL) as ws:
await ws.send(json.dumps(subscribe))
async def beat():
while True:
await asyncio.sleep(10)
await ws.send(json.dumps(heartbeat))
task = asyncio.create_task(beat())
try:
async for raw in ws:
msg = json.loads(raw)
if msg.get("cmd_id") == 22998:
yield msg["data"]
finally:
task.cancel()
except (websockets.ConnectionClosed, OSError):
await asyncio.sleep(3)
async def main():
api = AStockData()
code = "600519.SH"
bars = api.history(code, 500) # 先補歷史分鐘線
print("歷史K線", len(bars), "根,最後一根收盤價", bars[-1].close)
async for tick in api.stream_ticks([code, "000001.SZ"]):
print(tick["code"], tick["price"], tick["volume"])
asyncio.run(main())
開發效能飛躍:從底層除錯解脫,專注於市場 Alpha
將行情存取抽象為 history() 與 stream_ticks() 兩個核心標準方法後,團隊的開發模式將獲得實質改善:
- 回測與實盤無縫切換: 回測時直接以歷史資料庫模擬這兩個方法,實盤時無縫切入即時串流,策略核心邏輯一行都不必更動。
- 邊界狀態平滑對齊: 以分鐘起始時間戳作為唯一 Key,即時串流產生的動態資料即時覆蓋未完成的 REST 快照,徹底解決數據重複計算問題。
- 優化頻寬與 API 額度: 透過啟動時增量補齊機制配合本機快取,省去無謂的大量調用,規避頻率限制風險。
好的量化架構,就該把髒活累活擋在最底層。先把這套管線封裝好,未來無論策略如何迭代、資料源如何更換,你都能從容應對。
