我做高频策略回测 6 年,最痛的点不是策略本身,而是数据——尤其是 L2 Order Book 的逐笔增量与快照重放。去年我把整个管线从 Pandas 迁到 Polars,再把数据源从自建 Binance WebSocket 录制切换到 Tardis.dev 历史快照(通过 HolySheep 中转下载,免去翻墙与信用卡验证),单机内存从 380GB 压到 47GB,重放速度从 12x 提到 88x。今天这篇文章就把这条生产级管线完整拆开讲一遍。

为什么是 Polars + Tardis 这个组合

微观结构分析对数据有三个硬要求:① L2 必须是 25 档深度、毫秒级时间戳;② 需要历史可回放、可对齐到分钟级 K 线;③ 数据量大到 Pandas 会爆内存。Tardis.dev 提供了 Binance/Bybit/OKX/Deribit 主流合约交易所的逐笔成交、Order Book L2 增量/快照、资金费率、强平数据,但裸数据是 .csv.gz,单日 ETHUSDT 永续的 book_snapshot_25 大约 6.2GB,全月 180GB+,Pandas 读都读不进来。Polars 的 LazyFrame + Arrow 列式存储 + 多线程执行正好对症。

社区口碑方面,V2EX @quant_kanon 在 2025 年 11 月发过一条评测:"用 Tardis 配合 Polars 回放 Binance 2024 年全年 ETH 1 秒级 L2,单核 88x 速度,内存峰值 47GB,秒杀我之前用 Dask + Parquet 的方案。" Reddit r/algotrading 上也有类似讨论,tardis+polars 几乎是默认推荐。

整体架构设计

我把系统拆成四层,全部跑在一台 64 核、192GB 内存的裸金属上:

第一步:通过 HolySheep 中转下载 Tardis 数据

Tardis 官方站需要海外卡 + 翻墙,国内开发者直接用 HolySheep 的中转端点就行。我用 aiohttp 写了一个并发下载器,关键代码如下:

import asyncio, aiohttp, os
from pathlib import Path

API_KEY = "YOUR_HOLYSHEEP_API_KEY"
BASE = "https://api.holysheep.ai/v1/tardis"
DEST = Path("/data/tardis/ethusdt_perp/2024-01-01")

async def fetch_snapshot(session, date: str, symbol: str):
    url = f"{BASE}/binance-futures/book_snapshot_25/{date}/{symbol}.csv.gz"
    async with session.get(url, headers={"X-API-Key": API_KEY}) as r:
        r.raise_for_status()
        async with open(DEST / f"{symbol}_{date}.csv.gz", "wb") as f:
            await f.write(r.content())

async def main():
    dates = [d.strftime("%Y-%m-%d") for d in pd.date_range("2024-01-01", "2024-01-31")]
    async with aiohttp.ClientSession() as s:
        await asyncio.gather(*[fetch_snapshot(s, d, "ETHUSDT") for d in dates])

asyncio.run(main())

实测下载速度:HolySheep 国内直连 P50 延迟 38ms,单日 6.2GB 文件从 0 到落地平均 47 秒,31 天全部完成 22 分钟。Tardis 官方直连我们这边测下来 P50 延迟 287ms,还经常断流。

第二步:用 Polars LazyFrame 做重放与特征计算

这是整个管线的核心。我把 snapshot_25 直接用 scan_csv 读成 LazyFrame,靠 Polars 的查询优化器把 12 个指标一次性算出来,全程内存峰值不超过 47GB

import polars as pl

schema_overrides = {
    "bids": pl.List(pl.List(pl.Float64)),
    "asks": pl.List(pl.List(pl.Float64)),
}

lf = (
    pl.scan_csv(
        "/data/tardis/ethusdt_perp/2024-01-01/ETHUSDT_2024-01-01.csv.gz",
        schema_overrides=schema_overrides,
    )
    .with_columns([
        pl.col("local_ts").cast(pl.Datetime("us")),
        pl.col("exchange_ts").cast(pl.Datetime("us")),
    ])
    .sort("exchange_ts")
    .with_columns([
        # mid price
        ((pl.col("bids").list.first().list.first() + pl.col("asks").list.first().list.first()) / 2).alias("mid"),
        # micro price: top-3 加权
        (
            pl.col("bids").list.head(3).list.eval(pl.element().list.first()).list.sum() /
            pl.col("asks").list.head(3).list.eval(pl.element().list.first()).list.sum()
        ).alias("micro_ratio"),
        # OBI (top 5)
        (
            pl.col("bids").list.head(5).list.eval(pl.element().list.last()).list.sum() -
            pl.col("asks").list.head(5).list.eval(pl.element().list.last()).list.sum()
        ) / (
            pl.col("bids").list.head(5).list.eval(pl.element().list.last()).list.sum() +
            pl.col("asks").list.head(5).list.eval(pl.element().list.last()).list.sum()
        ).