【量化架構】如何擺脫數據對齊地獄?利用A股實時行情數據API打造統一進場管道

  • 5
  • 0

如果你曾負責過跨國資產配置或大中華區股票量化策略開發,一定碰過這種讓人抓狂的場景:回測時,你在本地硬碟讀取預先清洗好的分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() 兩個核心標準方法後,團隊的開發模式將獲得實質改善:

  1. 回測與實盤無縫切換: 回測時直接以歷史資料庫模擬這兩個方法,實盤時無縫切入即時串流,策略核心邏輯一行都不必更動。
  2. 邊界狀態平滑對齊: 以分鐘起始時間戳作為唯一 Key,即時串流產生的動態資料即時覆蓋未完成的 REST 快照,徹底解決數據重複計算問題。
  3. 優化頻寬與 API 額度: 透過啟動時增量補齊機制配合本機快取,省去無謂的大量調用,規避頻率限制風險。

好的量化架構,就該把髒活累活擋在最底層。先把這套管線封裝好,未來無論策略如何迭代、資料源如何更換,你都能從容應對。