很多量化项目的第一行代码是:
pd.read_csv("bars.csv")
问题通常也从这里开始。
这份数据属于哪个市场?时间戳是 bar 的开始还是结束?是否经过复权处理?最新一根是否已经收盘?volume=0 是无成交、停牌、休市,还是采集失败?这些问题如果在数据进入研究层之前没有答案,后面再精细的因子、参数搜索和 Walk-Forward ,都会建立在不稳定的输入上。
写清楚一份数据契约非常重要。
区分 DataFrame 与数据契约
DataFrame 是一个容器,不能回答“这行数据从哪里来”,也不能回答“当时策略究竟看到了什么”。数据契约则是行情接入层与研究层之间的一份明确约定:字段有什么含义、时间如何解释、哪些记录可以计算信号,以及出问题时如何追溯。
我会为每一根标准化 K 线至少保留以下字段。
| 字段 | 含义 | 为什么要保留 |
|---|---|---|
symbol |
供应商定义的标的代码 | 防止自定义 ticker 与真实数据代码混用。 |
bar_timestamp_utc |
原始时间戳转换后的 UTC 时间 | 让回测、定时任务和实时服务使用同一条时间轴。 |
interval |
例如 1m、5m、1d |
指标、下单频率和数据粒度必须显式匹配。 |
open/high/low/close |
标准化后的 OHLC | 不让策略层绑定某一家供应商的字段名。 |
volume/turnover |
成交量与成交额 | 为流动性分析与成本模型保留基础信息。 |
source |
数据来源标识 | 方便对账、迁移和多源比对。 |
source_trace |
一次获取操作的唯一标识 | 可以从研究结果倒查到具体请求。 |
raw_payload_sha256 |
原始响应的内容指纹 | 防止输入被静默替换而无人察觉。 |
ingested_at_utc |
系统接收数据的时间 | 区分市场事件时间与系统观察时间。 |
is_final |
是否允许进入“收盘后”策略 | 防止尚未完成的 bar 混进回测。 |
关键在于定义必须可以执行。比如,不能只写“所有数据有时间戳”,而要写:
入库时统一转换为 UTC ;凡是声明为收盘价策略的计算,只读取
is_final=True的记录。
研究层不应该认识供应商字段
一旦策略代码直接读取 close_price、trade_time 或某个供应商专有的符号格式,研究逻辑就和数据接入绑死了。之后替换数据源、复查某一次异常、增加第二份行情进行比对,成本都会迅速上升。
更理想的分层是这样的:
| 层级 | 责任 | 不应承担的责任 |
|---|---|---|
| 供应商适配层 | 发请求、鉴权、解析原始字段、保存原始响应 | 计算因子或决定交易。 |
| 标准化与质量层 | 转换字段、校验 OHLC 、标注数据状态、生成血缘信息 | 根据异常“猜测”如何修复。 |
| 研究层 | 读取稳定契约、构建特征、回测、评估稳健性 | 处理 token 、HTTP 超时、供应商 JSON 结构。 |
| 执行与风控层 | 处理事件时间、订单状态、成本与限额 | 重写历史数据或掩盖数据缺陷。 |
这样做的直接收益,是策略代码不必在意行情来自 HTTP 、WebSocket 、数据库还是 CSV 。研究层只依赖契约;接入层替换时,最多改适配器与测试,而不是重写全部实验。
可审计适配器小 Demo
下面的代码以 AllTick 的单标的 K 线 REST 接口作为具体适配器示例。
# pip install requests pandas
from __future__ import annotations
import hashlib
import json
import os
from datetime import datetime, timezone
from pathlib import Path
from uuid import uuid4
import pandas as pd
import requests
# This URL is one provider-specific adapter detail, not a research-layer dependency.
KLINE_URL = "https://quote.alltick.co/quote-stock-b-api/kline"
RAW_DIR = Path("data/raw/market_data")
def fetch_bars(code: str, *, interval_code: int = 1, count: int = 200) -> tuple[dict, str]:
"""Fetch raw vendor data and return (payload, trace_id)."""
token = os.getenv("ALLTICK_TOKEN")
if not token:
raise RuntimeError("Set ALLTICK_TOKEN in the environment; never commit credentials.")
if not 1 <= count <= 500:
raise ValueError("This adapter accepts 1–500 bars per historical request.")
trace = str(uuid4())
query = {
"trace": trace,
"data": {
"code": code,
"kline_type": interval_code,
"kline_timestamp_end": 0,
"query_kline_num": count,
"adjust_type": 0,
},
}
response = requests.get(
KLINE_URL,
params={"token": token, "query": json.dumps(query, separators=(",", ":"))},
timeout=10,
)
response.raise_for_status()
payload = response.json()
if payload.get("ret") != 200:
raise RuntimeError(f"Market-data application error: {payload}")
if payload.get("trace") != trace:
raise RuntimeError("Response trace mismatch; do not ingest an uncorrelated response.")
return payload, trace
def persist_raw(payload: dict, trace: str) -> tuple[Path, str]:
"""Save the source record before any interpretation occurs."""
RAW_DIR.mkdir(parents=True, exist_ok=True)
raw = json.dumps(payload, ensure_ascii=False, sort_keys=True).encode("utf-8")
digest = hashlib.sha256(raw).hexdigest()
path = RAW_DIR / f"{trace}.json"
path.write_bytes(raw)
return path, digest
def to_contract(payload: dict, trace: str, raw_sha256: str) -> pd.DataFrame:
"""Map vendor fields to the vendor-neutral research contract."""
data = payload["data"]
raw_bars = pd.DataFrame(data["kline_list"])
for field in ["open_price", "high_price", "low_price", "close_price", "volume", "turnover"]:
raw_bars[field] = pd.to_numeric(raw_bars[field], errors="coerce")
bars = pd.DataFrame({
"symbol": data["code"],
"bar_timestamp_utc": pd.to_datetime(raw_bars["timestamp"].astype("int64"), unit="s", utc=True),
"interval": "1m",
"open": raw_bars["open_price"],
"high": raw_bars["high_price"],
"low": raw_bars["low_price"],
"close": raw_bars["close_price"],
"volume": raw_bars["volume"],
"turnover": raw_bars["turnover"],
"source": "market_data_provider",
"source_trace": trace,
"raw_payload_sha256": raw_sha256,
"ingested_at_utc": datetime.now(timezone.utc).isoformat(),
}).sort_values("bar_timestamp_utc").reset_index(drop=True)
# Conservative default: the newest observation may be a currently forming bar.
bars["is_final"] = True
if not bars.empty:
bars.loc[bars.index[-1], "is_final"] = False
return bars
payload, trace = fetch_bars("AAPL.US", interval_code=1, count=200)
raw_path, raw_hash = persist_raw(payload, trace)
bars = to_contract(payload, trace, raw_hash)
print(f"raw source record: {raw_path}")
print(bars.tail(3).to_string(index=False))
有三个要点:
第一,原始响应先落盘,再转换字段。原始 JSON 是供应商给出的事实记录;标准化 K 线是我们的解释。两份都保留,未来才能在解析规则变化、发现时间语义错误或需要复现实验时重新生成研究表。
第二,对每次输入建立血缘。trace 将一次数据获取与输出关联起来,哈希记录原始内容。将来某个异常信号出现时,应该可以追问:“这根 K 线来自哪次获取?原始响应还在吗?当时转换器的版本是什么?”
第三,把最新 bar 作为暂定记录处理。这一点不依赖某一家供应商。实时数据中的最新 bar 常常仍在形成,如果策略声称“仅使用已收盘数据”,就不能让它直接参与信号。供应商文档若明确说明最新 bar 的语义,应据此实现更精确的最终状态判定。
校验的职责是隔离
成功拿到 HTTP200,不代表数据已经适合回测。接入层至少应做一些低成本、强约束的检查。
| 检查 | 失败后的合理做法 | 为什么 |
|---|---|---|
| 请求代码与返回代码一致 | 拒绝入库并告警 | 标的映射错误会污染全部下游结论。 |
| 时间戳递增且无重复 | 进入隔离表 | 重复 bar 会扭曲收益与成交次数。 |
low <= min(open, close) <= max(open, close) <= high |
保留原始数据,等待定位 | OHLC 不自洽通常不应被静默修复。 |
| OHLC 为正,成交量与成交额非负 | 隔离 | fillna(0) 会把质量问题伪装成市场状态。 |
| 最新 bar 的状态可解释 | 默认不纳入收盘价策略 | 防止未来函数。 |
| 原始文件和版本可定位 | 停止进入研究库 | 无法追溯的数据不应支持结论。 |
一个简单的校验器就足够成为第一道防线:
def validate_ohlcv(frame: pd.DataFrame) -> list[str]:
errors: list[str] = []
if frame["bar_timestamp_utc"].duplicated().any():
errors.append("duplicate timestamps")
if not frame["bar_timestamp_utc"].is_monotonic_increasing:
errors.append("timestamps are not sorted")
prices = frame[["open", "high", "low", "close"]]
if prices.isna().any().any() or (prices <= 0).any().any():
errors.append("non-positive or missing OHLC")
ohlc_ok = (
(frame["low"] <= frame[["open", "close"]].min(axis=1))
& (frame[["open", "close"]].max(axis=1) <= frame["high"])
)
if not ohlc_ok.all():
errors.append("inconsistent OHLC relationship")
if (frame[["volume", "turnover"]] < 0).any().any():
errors.append("negative volume or turnover")
return errors
issues = validate_ohlcv(bars)
if issues:
raise ValueError(f"Quarantine this batch; do not backtest it: {issues}")
research_bars = bars.loc[bars["is_final"]].copy()
这里的核心原则是:不要让异常安静地变成正常值。
例如,缺失 bar 可能对应休市、停牌、无成交、供应商延迟或自己的采集故障。直接前向填充上一根 close ,确实能让均线继续算下去,但也会改变波动率、成交频率和成本估计。先确认缺失属于什么状态,再定义研究层应如何处理。
历史初始化与在线刷新要分开
任何可靠的研究库,都应把“补历史”和“更新最新数据”视作两类任务。
| 任务 | 目标 | 推荐结果 |
|---|---|---|
| 历史初始化 | 获取并冻结一个明确的数据快照 | 本地版本化历史表,附带来源、时间和哈希。 |
| 增量刷新 | 更新最近发生变化的 bar | 幂等 upsert ,并保留获取时间。 |
| 实时消费 | 接收并持久化市场事件 | 有去重键、重连策略和背压控制的事件流。 |
策略应该只看到稳定的输入
当接入层完成后,策略代码只应依赖标准契约,而不是任何供应商的 JSON 、鉴权逻辑或字段命名。
# Strategy code consumes a contract, not a vendor response.
signal_input = research_bars.loc[
research_bars["is_final"]
& (research_bars["interval"] == "1m")
].copy()
signal_input["ret_1"] = signal_input["close"].pct_change()
signal_input["ma_20"] = signal_input["close"].rolling(20, min_periods=20).mean()
这段代码不会自动产生 alpha ,但它让我们能清楚地回答:一个回测结果来自哪一版数据、当时可用的信息是什么、是否混入未收盘 bar 、以及几个月后能否被另一个人重新跑出来。
先把数据边界做扎实,后面的交易成本、滑点、Walk-Forward 、状态机和实时系统,才有可讨论的共同语言。