首页 / 文章 / 构建市场时光机:使用 Python 和 WebSocket 重放交易会话
← 返回
IT技术

构建市场时光机:使用 Python 和 WebSocket 重放交易会话

✍️ zhirenhun 📅 2026/8/30 👁 124 阅读 ⏱ 163 分钟
构建市场时光机:使用 Python 和 WebSocket 重放交易会话

历史市场数据通常以完成的数据集形式到达。这对分析很方便,但与交易软件体验实时市场的方式截然不同。在生产环境中,事件会逐一到达,未来是未知的,每个决策仅依赖于迄今为止发生的事情。

在本教程中,我们将使用历史逐笔数据重建这种体验。我们将从 EODHD 获取完整的 AAPL 交易时段,将超过一百万笔交易归一化为确定性事件带,并通过可控的市场时钟按照原始时间顺序重播它们。

在此过程中,我们将添加可调的播放速度、暂停和恢复控制、搜索功能,以及一个通过 REST 暴露控制并通过 WebSocket 流式传输交易的 FastAPI 服务。

我们还将构建一个独立的消费者,仅根据其接收到的事件计算滚动 VWAP 和市场状态。完成后,我们将拥有一个完整的本地回放系统,能够将已经结束的交易日作为定时流喂送给事件驱动的软件,在搜索后正确重建下游状态,并通过自动化测试验证结果。

目录

先决条件

开始之前,请确保您已具备:

本教程中的 shell 命令使用 Unix 风格语法,因此可以直接在 macOS 和 Linux 上运行。在 Windows 上,您可以通过 WSL、Git Bash 运行它们,或使用等效的 PowerShell 命令。

我们正在构建的内容

在触碰代码之前,先查看完整系统会有所帮助。回放引擎将从 EODHD 获取历史交易,将其转换为一致的内部格式,恢复其时间,并像交易日正在重演一样将其流式传输给独立的消费者。

完整流程如下:

每一层只负责一项工作。加载器负责检索并保存原始的历史会话。规范化器验证这些记录并将其转换为确定性的回放磁带。时钟将历史时间戳映射到墙上时间,而回放会话则添加了诸如开始、暂停、恢复、速度调节、查找和停止等控制。

FastAPI 围绕这个回放引擎运行。REST 端点构成控制平面,而 WebSocket 负责传输实际的交易和回放控制事件。在另一侧,消费者仅通过该流获得的数据维护自己的滚动状态。

我们将在项目结构中保持这些职责的分离:

market-time-machine/
├── data/
│   ├── raw/
│   └── processed/
├── replay/
│   ├── __init__.py
│   ├── config.py
│   ├── loader.py
│   ├── events.py
│   ├── clock.py
│   └── session.py
├── api/
│   ├── __init__.py
│   ├── server.py
│   └── run.py
├── consumer/
│   ├── __init__.py
│   └── consumer.py
├── tests/
│   ├── __init__.py
│   ├── conftest.py
│   └── test_replay.py
├── .env
├── .gitignore
└── pytest.ini

整个构建的关键规则很简单:消费者只能知道回放流中已经到达的内容,永远不应该从历史磁带读取超前数据。正是这个约束,使得时序、暂停/恢复行为以及查找后的状态重构值得正确实现。

设置 Python 项目

首先创建项目目录并安装我们将用于数据检索、回放时序、API 层、WebSocket 通信和测试的软件包。

mkdir -p market-time-machine/data/raw
mkdir -p market-time-machine/data/processed
mkdir -p market-time-machine/replay
mkdir -p market-time-machine/api
mkdir -p market-time-machine/consumer
mkdir -p market-time-machine/tests

cd market-time-machine

pip install requests fastapi "uvicorn[standard]" websockets httpx python-dotenv numpy pytest pytest-asyncio

replayapiconsumertests 目录中创建空的 __init__.py 文件,使 Python 将每个目录视为一个包:

replay/__init__.py
api/__init__.py
consumer/__init__.py
tests/__init__.py

我们将从 来自 EODHD 的历史交易 获取数据,因此在项目根目录创建一个 .env 文件,并将您的 API 密钥存放其中:

下载的会话也会相当大,因此凭证和本地市场数据文件都不应提交。请创建 .gitignore

.env
data/
__pycache__/
*.pyc
.ipynb_checkpoints/

注意: 如果您没有 EODHD API 密钥,可以通过 打开 EODHD 开发者账户 轻松获取。

在此设置之后,项目应如下所示:

market-time-machine/
├── data/
│   ├── raw/
│   └── processed/
├── replay/
│   └── __init__.py
├── api/
│   └── __init__.py
├── consumer/
│   └── __init__.py
├── tests/
│   └── __init__.py
├── .env
└── .gitignore

raw/ 目录将保存从 EODHD 收到的响应,而 processed/ 目录则保存我们基于这些响应构建的归一化回放磁带。

从 EODHD 下载完整交易会话

回放引擎在能够恢复时间概念之前,需要一个完整的交易会话。我们将使用 EODHD 的历史 tick API 检索 2026 年 7 月 15 日的 AAPL 成交数据,但会将检索层与回放相关的所有内容保持分离。

以下两个文件负责项目的这一部分:

market-time-machine/
└── replay/
    ├── __init__.py
    ├── config.py
    └── loader.py

config.py 将共享的 API、路径和市场会话设置集中在一个地方。loader.py 使用这些设置来检索会话,并将原始响应保存在 data/raw/ 目录下。

创建 replay/config.py

添加以下内容:

import os
from pathlib import Path
from dotenv import load_dotenv

ROOT = Path(__file__).resolve().parent.parent
load_dotenv(ROOT / ".env")

TOKEN = os.environ.get("EODHD_API_TOKEN")
TICKS_URL = "https://eodhd.com/api/ticks/"

RAW = ROOT / "data" / "raw"
PROCESSED = ROOT / "data" / "processed"

MARKET_TZ = "America/New_York"
OPEN = "09:30:00"
CLOSE = "16:00:00"

MAX_LIMIT = 10_000
MIN_WINDOW_S = 1
CLOSE_GRACE_S = 5

FIELDS = ("mkt", "price", "seq", "shares", "sl", "sub_mkt", "ts")
NON_LAST_SALE = frozenset("IWVT47")

def token():
    if not TOKEN:
        raise RuntimeError("EODHD_API_TOKEN not set")
    return TOKEN

def redact(text):
    return str(text).replace(TOKEN, "") if TOKEN else str(text)

美国股票常规交易时段定义在 America/New_York,而不是使用固定的 UTC 时间戳。这一点很重要,因为 09:30 对应的 UTC 时间会随夏令时而变化。

我们还通过 CLOSE_GRACE_S 将请求窗口延长至 16:00 之后五秒。本教程使用的会话在 16:00:00 之后立即包含收盘活动,因此宽限窗口能够将这些记录保留在下载中。

创建 replay/loader.py

单次大规模请求不是获取密集 tick 数据会话的安全方式。全天活动变化显著,任何达到配置的 10,000 条记录上限的请求都可能代表被截断的区间。

相反,加载器会根据前一次响应的数据密度来调整其请求窗口。

创建 replay/loader.py

import json, time
from datetime import datetime
from zoneinfo import ZoneInfo

import requests

from . import config

def fetch(symbol, frm, to, limit=None):
    limit = limit or config.MAX_LIMIT

    r = requests.get(config.TICKS_URL, timeout=180, params={
        "s": symbol,
        "from": frm,
        "to": to,
        "limit": limit,
        "api_token": config.token(),
        "fmt": "json"
    })

    if r.status_code != 200:
        raise RuntimeError(
            f"HTTP {r.status_code} {config.redact(r.text[:200])}"
        )

    return r.json()

def bounds(date_str, grace=None):
    grace = config.CLOSE_GRACE_S if grace is None else grace
    tz = ZoneInfo(config.MARKET_TZ)
    d = datetime.strptime(date_str, "%Y-%m-%d").date()

    def at(hms):
        h, m, s = map(int, hms.split(":"))
        return datetime(
            d.year, d.month, d.day, h, m, s, tzinfo=tz
        ).timestamp()

    return int(at(config.OPEN)), int(at(config.CLOSE)) + grace

def fetch_session(symbol, date_str, tag="session", window=None,
                  force=False, verbose=True):

    raw = config.RAW / f"{symbol}_{date_str}_{tag}.jsonl"
    man = config.RAW / f"{symbol}_{date_str}_{tag}.manifest.json"

    if raw.exists() and man.exists() and not force:
        m = json.loads(man.read_text())
        print(f"cached {raw.name}: {m['ticks']:,} ticks")
        return m, raw

    start, end = window or bounds(date_str)
    cursor, win = start, 30

    total = pages = retries = 0
    first_ts = last_ts = None
    seen_fields = set()
    t0 = time.perf_counter()

    with open(raw, "w") as fh:
        while cursor < end:
            b = min(cursor + win, end)
            span = b - cursor

            payload = fetch(symbol, cursor, b)
            n = len(payload.get("ts", []))

            if n >= config.MAX_LIMIT:
                if span <= config.MIN_WINDOW_S:
                    raise RuntimeError(
                        f"second {cursor} has >= {config.MAX_LIMIT} ticks "
                        "and cannot be paginated"
                    )

                win = max(1, span // 2)
                retries += 1
                continue

            if n:
                seen_fields.update(payload.keys())

                if first_ts is None:
                    first_ts = payload["ts"][0]

                last_ts = payload["ts"][-1]

                fh.write(json.dumps({
                    "from": cursor,
                    "to": b,
                    "payload": payload
                }) + "\n")

            total += n
            pages += 1
            cursor = b

            density = n / span if span else 0
            win = int(min(
                1800,
                max(1, config.MAX_LIMIT * 0.75 / max(density, 0.01))
            ))

            if verbose and pages % 20 == 0:
                pct = 100 * (cursor - start) / (end - start)
                print(f"{pct:5.1f}% {total:,} ticks")

    m = {
        "symbol": symbol,
        "date": date_str,
        "tag": tag,
        "ticks": total,
        "pages": pages,
        "retries": retries,
        "window_from_utc": start,
        "window_to_utc": end,
        "first_timestamp_ms": first_ts,
        "last_timestamp_ms": last_ts,
        "fields": sorted(seen_fields),
        "elapsed_s": round(time.perf_counter() - t0, 1),
        "api_calls": pages * 10
    }

    man.write_text(json.dumps(m, indent=2))
    return m, raw

def read_pages(path):
    with open(path) as fh:
        for line in fh:
            if line.strip():
                yield json.loads(line)

加载器初始窗口为30秒。如果该间隔达到记录上限,则会改用更小的窗口重试,而不是接受可能不完整的响应。在较为平静的时期,后续窗口可扩展至最长30分钟。

每个被接受的响应会在任何归一化之前直接写入JSONL。一个清单会随其一起存储,包含会话边界、滴答计数、时间戳、观察到的字段以及检索统计信息。

现在获取完整的AAPL会话:

from replay.loader import fetch_session

SYMBOL = "AAPL"
DATE = "2026-07-15"

print("=== fullday ===")

manifest, raw_path = fetch_session(
    SYMBOL,
    DATE,
    tag="fullday"
)

print(
    f" window {manifest['window_from_utc']}..{manifest['window_to_utc']} | "
    f"{manifest['ticks']:,} ticks, {manifest['pages']} pages, "
    f"{manifest['retries']} retries | "
    f"{manifest['elapsed_s']}s, {manifest['api_calls']} metered api calls"
)

print(
    f" first_ts {manifest['first_timestamp_ms']} "
    f"last_ts {manifest['last_timestamp_ms']}"
)

print(" fields:", manifest["fields"])

最终的干净运行复用了已下载的会话,产生了:

我们现在有 1,032,411 条原始交易记录,覆盖完整的常规交易时段和收盘宽限窗口。145 个被接受的页面和 16 次重试也说明,对如此密集的 tick 数据假设固定请求窗口并不稳健。

这些记录仍然保存为从 EODHD 获得的原始样子。但在回放引擎能够使用它们之前,它们需要转换为确定的内部事件序列。

将 tick 数据归一化为回放带

加载器为我们提供了完整的会话,但回放引擎不应直接使用 EODHD 的原始响应格式。tick 接口返回的字段(如时间戳、价格、规模、序列号和市场代码)是以并行数组形式给出的。

在回放之前,我们需要核对这些数组是否对齐,建立确定的事件顺序,去除重复项,并将结果转换为一种内部格式。

该逻辑位于 replay/events.py 中:

market-time-machine/
└── replay/
    ├── config.py
    ├── loader.py
    └── events.py

我们这里会使用两个对象。TradeEvent 表示将最终通过 WebSocket 传输的单笔交易。TradeTape 以列式 NumPy 数组高效存储完整会话,并且仅在需要时才实体化单个 TradeEvent 对象。

创建 replay/events.py

创建 replay/events.py,内容如下:

from dataclasses import dataclass
import numpy as np

from . import config
from .loader import read_pages

@dataclass(frozen=True)
class TradeEvent:
    symbol: str
    timestamp_ms: int
    price: float
    size: int
    sequence: int
    market: str
    sub_market: str
    sale_condition: str
    source: str = "replay"

    def to_wire(self):
        sl = self.sale_condition

        return {
            "type": "trade",
            "symbol": self.symbol,
            "timestamp_ms": self.timestamp_ms,
            "price": self.price,
            "size": self.size,
            "sequence": self.sequence,
            "source": self.source,
            "metadata": {
                "market": self.market,
                "sub_market": self.sub_market or None,
                "sale_condition": sl,
                "odd_lot": "I" in sl,
                "zero_size": self.size == 0,
                "last_sale_eligible": not (
                    set(sl) & config.NON_LAST_SALE
                )
            }
        }


class TradeTape:
    def __init__(self, symbol, ts, price, size, seq, mkt, sub, sl):
        self.symbol = symbol
        self.ts = ts
        self.price = price
        self.size = size
        self.seq = seq
        self.mkt = mkt
        self.sub = sub
        self.sl = sl

    def __len__(self):
        return len(self.ts)

    def __getitem__(self, i):
        return TradeEvent(
            self.symbol,
            int(self.ts[i]),
            float(self.price[i]),
            int(self.size[i]),
            int(self.seq[i]),
            str(self.mkt[i]),
            str(self.sub[i]),
            str(self.sl[i])
        )

    def index_at(self, ts_ms):
        return int(np.searchsorted(self.ts, ts_ms, side="left"))

    def span(self):
        if not len(self):
            return None, None

        return int(self.ts[0]), int(self.ts[-1])

    def save(self, path):
        np.savez_compressed(
            path,
            ts=self.ts,
            price=self.price,
            size=self.size,
            seq=self.seq,
            mkt=self.mkt,
            sub=self.sub,
            sl=self.sl,
            symbol=np.array([self.symbol])
        )

    @classmethod
    def load(cls, path):
        z = np.load(path, allow_pickle=False)

        return cls(
            str(z["symbol"][0]),
            z["ts"],
            z["price"],
            z["size"],
            z["seq"],
            z["mkt"],
            z["sub"],
            z["sl"]
        )


def normalize(raw_path, symbol, verbose=True):
    cols = {k: [] for k in config.FIELDS}
    pages = 0

    for page in read_pages(raw_path):
        pages += 1
        p = page["payload"]

        lens = {k: len(p.get(k, [])) for k in config.FIELDS}

        if len(set(lens.values())) != 1:
            raise ValueError(
                f"ragged page {page['from']}: {lens}"
            )

        for k in config.FIELDS:
            cols[k].extend(p[k])

    ts = np.asarray(cols["ts"], dtype=np.int64)
    price = np.asarray(cols["price"], dtype=np.float64)
    size = np.asarray(cols["shares"], dtype=np.int64)
    seq = np.asarray(cols["seq"], dtype=np.int64)
    mkt = np.asarray(cols["mkt"], dtype=str)
    sub = np.asarray(cols["sub_mkt"], dtype=str)
    sl = np.asarray(cols["sl"], dtype=str)

    raw_n = len(ts)

    def arrays(mask):
        return tuple(
            a[mask]
            for a in (ts, price, size, seq, mkt, sub, sl)
        )

    keep = (
        (ts > 0)
        & np.isfinite(price)
        & (price > 0)
        & (size >= 0)
    )

    ts, price, size, seq, mkt, sub, sl = arrays(keep)

    order = np.lexsort((seq, ts))
    ts, price, size, seq, mkt, sub, sl = arrays(order)

    dup = np.zeros(len(ts), dtype=bool)

    if len(ts) > 1:
        dup[1:] = (
            (ts[1:] == ts[:-1])
            & (seq[1:] == seq[:-1])
        )

    ts, price, size, seq, mkt, sub, sl = arrays(~dup)

    tape = TradeTape(
        symbol,
        ts,
        price,
        size,
        seq,
        mkt,
        sub,
        sl
    )

    odd = sum("I" in str(s) for s in sl)
    elig = sum(
        not (set(str(s)) & config.NON_LAST_SALE)
        for s in sl
    )

    rep = {
        "pages": pages,
        "raw": raw_n,
        "kept": len(ts),
        "dropped": raw_n - len(ts) - int(dup.sum()),
        "dupes": int(dup.sum()),
        "zero_size": int((size == 0).sum()),
        "odd_lot": int(odd),
        "last_sale_eligible": int(elig),
        "seq_strict": bool(
            np.all(seq[1:] > seq[:-1])
        ) if len(seq) > 1 else True,
        "span": tape.span()
    }

    if verbose:
        n = max(1, len(ts))

        print(
            f"{rep['raw']:,} raw -> {rep['kept']:,} kept "
            f"({rep['dupes']} dupes, {rep['dropped']} invalid)"
        )

        print(
            f"zero-size {100*rep['zero_size']/n:.1f}% | "
            f"odd-lot {100*odd/n:.1f}% | "
            f"last-sale-eligible {100*elig/n:.1f}%"
        )

        print(
            f"seq strictly increasing: {rep['seq_strict']}"
        )

    return tape, rep

第一次验证在构建任何交易之前进行。由于源字段以并行数组的形式到达,页面上的每个字段必须包含相同数量的观测值;否则,合并它们可能会悄悄地把一笔交易的价格附到另一笔交易的时间戳上。

随后,这些数组被转换为 NumPy,移除基本的无效记录,并按照 (timestamp, sequence) 对交易进行排序。时间戳提供时间顺序,而序列号在多笔交易共享同一毫秒时提供确定性的顺序。

随后会删除时间戳和序列完全相同的精确重复记录。TradeTape 将结果列保留为数组,而不是分配超过一百万个永久的 Python 对象,因而全天会话在内存占用上显著降低。

现在对原始会话进行归一化,并将其保存到 data/processed/ 目录下:

import json

from replay import config
from replay.events import normalize

tape, report = normalize(raw_path, SYMBOL)

tape.save(
    config.PROCESSED / f"{SYMBOL}_{DATE}_fullday.npz"
)

lo, hi = tape.span()

print(
    f"span {lo}..{hi} "
    f"({(hi-lo)/3_600_000:.2f} market hours)"
)

print("first 3 normalized events:")

for i in range(3):
    print(json.dumps(tape[i].to_wire()))

实际的归一化运行产生了:

在超过一百万条原始观测中,只有两条重复记录被剔除,且没有任何记录未能通过基本的时间戳、价格或大小检查。更重要的是,对于回放,归一化后的序列是严格递增的。

创建用于基准和测试的较小带

全天带将驱动最终的回放。不过,对于时序基准测试和自动化测试,我们不必每次都跑完全部 6.5 小时。

我们将直接从归一化的全天带中提取东部时间 12:00 至 12:15 的 15 分钟切片:

from datetime import datetime
from zoneinfo import ZoneInfo

from replay.events import TradeTape

tz = ZoneInfo(config.MARKET_TZ)

quiet_start = int(
    datetime(
        2026, 7, 15, 12, 0,
        tzinfo=tz
    ).timestamp() * 1000
)

quiet_end = quiet_start + 15 * 60_000

i = tape.index_at(quiet_start)
j = tape.index_at(quiet_end)

quiet_tape = TradeTape(
    tape.symbol,
    tape.ts[i:j],
    tape.price[i:j],
    tape.size[i:j],
    tape.seq[i:j],
    tape.mkt[i:j],
    tape.sub[i:j],
    tape.sl[i:j]
)

quiet_tape.save(config.PROCESSED / f"{SYMBOL}_{DATE}_quiet15m.npz")

我们现在有两段已处理的磁带:一段用于端到端回放的完整会话,另一段是用于可重复计时和控制测试的较小真实市场区间。

构建历史回放时钟

我们现在拥有一个确定的交易序列,但仍缺少让这些交易像市场流一样运行的机制。直接遍历磁带时,Python 会以机器允许的最快速度处理会话。

回放时钟通过把历史市场时间映射到真实墙上时间来解决这个问题。它还能在不改变原始时间戳的情况下调节播放速度。

一种简单的做法可能是:在每对交易之间的历史间隔上让程序休眠。

gap = (next_ts - current_ts) / 1000
await asyncio.sleep(gap / speed)

10x 时,500 ms 的历史间隔变为 50 ms。在 100x 时,它变为 5 ms。

问题在于 asyncio.sleep() 仅保证执行将在请求的延迟之后恢复。如果每次睡眠稍微延迟唤醒,而下一次延迟是从该延迟唤醒点测量的,这些误差可能在长时间的回放中累积。

相反,我们将整个回放锚定到 time.monotonic()

historical elapsed time
        ÷
replay speed
        +
wall-clock start
        =
target wall-clock time

因此,所有事件均相对于同一锚点调度,而不以前一个事件结束时间为参照。

创建 replay/clock.py

将时钟添加到 replay 包中:

market-time-machine/
└── replay/
    ├── config.py
    ├── loader.py
    ├── events.py
    └── clock.py

创建 replay/clock.py

import asyncio, time
import numpy as np

MIN_SLEEP_S = 0.0005

class ReplayClock:
    def __init__(self, start_ms, speed=1.0):
        self.speed = float(speed)
        self._anchor_ms = float(start_ms)
        self._anchor_wall = None
        self.running = False
        self.epoch = 0

    def start(self):
        self._anchor_wall = time.monotonic()
        self.running = True
        return self

    def now_ms(self, now=None):
        if not self.running or self._anchor_wall is None:
            return self._anchor_ms

        now = now if now is not None else time.monotonic()

        return (
            self._anchor_ms
            + (now - self._anchor_wall) * 1000 * self.speed
        )

    def wall_for(self, ms):
        return (
            self._anchor_wall
            + (ms - self._anchor_ms) / 1000 / self.speed
        )

    def _reanchor(self, ms):
        self._anchor_ms = float(ms)
        self._anchor_wall = time.monotonic()
        self.epoch += 1

    def set_speed(self, speed):
        self._reanchor(self.now_ms())
        self.speed = float(speed)

    def pause(self):
        if self.running:
            self._anchor_ms = self.now_ms()
            self.running = False

    def resume(self):
        if not self.running:
            self._anchor_wall = time.monotonic()
            self.running = True
            self.epoch += 1

    def seek(self, ms):
        self._reanchor(ms)


def new_stats(speed):
    return {
        "emitted": 0,
        "batches": 0,
        "lateness": [],
        "dropped": 0,
        "speed": speed,
        "wall0": None,
        "market0": None,
        "market1": None
    }


def summarize(st):
    if not st["lateness"]:
        return {
            "emitted": st["emitted"],
            "batches": st["batches"]
        }

    a = np.asarray(st["lateness"])

    wall = (
        time.monotonic() - st["wall0"]
        if st["wall0"] else 0.0
    )

    mkt = (
        (st["market1"] - st["market0"]) / 1000
        if st["market0"] is not None else 0.0
    )

    ok = st["dropped"] == 0 and wall > 0
    realized = round(mkt / wall, 2) if ok else None

    return {
        "emitted": st["emitted"],
        "batches": st["batches"],
        "mean_batch": round(
            st["emitted"] / max(1, st["batches"]), 1
        ),
        "market_s": round(mkt, 3),
        "wall_s": round(wall, 3),
        "requested_speed": st["speed"],
        "realized_speed": realized,
        "speed_error_pct": (
            round(
                100 * (realized - st["speed"]) / st["speed"],
                2
            )
            if ok else None
        ),
        "lateness_p50_ms": round(
            float(np.percentile(a, 50)), 2
        ),
        "lateness_p95_ms": round(
            float(np.percentile(a, 95)), 2
        ),
        "lateness_max_ms": round(
            float(a.max()), 2
        ),
        "reanchor_batches_dropped": st["dropped"]
    }


async def replay_batches(
    tape,
    clock,
    start,
    stats,
    max_batch=4096
):
    i, n = start, len(tape)
    last_epoch = clock.epoch

    if stats["wall0"] is None:
        stats["wall0"] = time.monotonic()
        stats["market0"] = int(tape.ts[start])

    while i < n:
        if not clock.running:
            await asyncio.sleep(0.005)
            continue

        now = time.monotonic()

        j = min(
            int(
                np.searchsorted(
                    tape.ts,
                    clock.now_ms(now),
                    side="right"
                )
            ),
            n,
            i + max_batch
        )

        if j > i and not clock.running:
            continue

        if j > i:
            if clock.epoch == last_epoch:
                targets = clock.wall_for(
                    tape.ts[i:j].astype(np.float64)
                )

                stats["lateness"].extend(
                    ((now - targets) * 1000).tolist()
                )
            else:
                stats["dropped"] += 1
                last_epoch = clock.epoch

            stats["emitted"] += j - i
            stats["batches"] += 1
            stats["market1"] = int(tape.ts[j - 1])

            yield i, j
            i = j
            continue

        wait = clock.wall_for(float(tape.ts[i])) - now

        await asyncio.sleep(
            wait if wait > MIN_SLEEP_S else 0
        )

now_ms() 告诉我们回放目前在历史市场时间的位置。wall_for() 执行相反的转换,并告诉我们历史时间戳应在机器的单调时钟上何时到期。

暂停、恢复、调速和查找可以随后重新锚定该映射,而无需修改底层磁带。

另一个重要部分是批处理。在高回放速度下,为每笔交易安排一次 sleep 会产生可观的开销。replay_batches() 则会询问市场时间已经推进了多少,并释放所有已到期的交易,直至达到配置的批次大小。

如果事件循环略有滞后,下一批次会变大,而不是引入另一个人为延迟。

基准测试回放时钟

现在加载我们在上一节创建的午间磁带,测试前 30 秒的市场时间:

import numpy as np

from replay import config
from replay.events import TradeTape
from replay.clock import (
    ReplayClock,
    replay_batches,
    new_stats,
    summarize
)

tape = TradeTape.load(
    config.PROCESSED / "AAPL_2026-07-15_quiet15m.npz"
)

end = int(
    np.searchsorted(
        tape.ts,
        tape.ts[0] + 30_000,
        side="right"
    )
)

print(
    f"{end:,} events in the first "
    "30 market seconds of AAPL quiet15m\n"
)

async def measure():
    print(
        f"{'speed':>6} {'market_s':>9} "
        f"{'wall_s':>8} {'realized':>9} "
        f"{'err_%':>7} {'p50_ms':>7} "
        f"{'p95_ms':>7} {'max_ms':>7}"
    )

    for speed in [1, 10, 50, 100]:
        clock = ReplayClock(
            tape.ts[0],
            speed
        ).start()

        st = new_stats(speed)

        async for i, j in replay_batches(
            tape,
            clock,
            0,
            st
        ):
            if j >= end:
                break

        r = summarize(st)

        print(
            f"{r['requested_speed']:>6} "
            f"{r['market_s']:>9} "
            f"{r['wall_s']:>8} "
            f"{r['realized_speed']:>9} "
            f"{r['speed_error_pct']:>7} "
            f"{r['lateness_p50_ms']:>7} "
            f"{r['lateness_p95_ms']:>7} "
            f"{r['lateness_max_ms']:>7}"
        )

await measure()

实际运行产生了:

历史市场时间的三十秒在 1x 倍速下用了 30.001 秒,在 10x 倍速下用了 3.001 秒,在 50x 倍速下用了 0.6 秒,在 100x 倍速下用了 0.3 秒。实现的速度因此与我们请求的非常接近。

延迟值表明调度器错过各个事件截止时间的程度。例如,在 10x 倍速下,中位延迟为 0.36 ms,第 95 百分位为 1.12 ms,而此次运行中观察到的最坏情况为 11.75 ms

这些数字衡量的是回放时钟本身。它们不是端到端的 WebSocket 延迟测量,这仍然是基于 Python 事件循环的尽力调度,而非交易所级别的时序。

为回放会话添加播放控制

回放时钟知道交易何时到期,但它不知道回放目前的位置,也不清楚是否应该进行播放。我们需要另一层来负责磁带、跟踪当前光标、管理事件队列,并协调诸如开始、暂停、恢复、速度变化、查找和停止等控制。

该逻辑位于 replay/session.py 中:

market-time-machine/
└── replay/
    ├── config.py
    ├── loader.py
    ├── events.py
    ├── clock.py
    └── session.py

保持这一区分很有用:时钟负责时间,而会话负责状态。

回放会话会经过一组有限的状态:

创建 replay/session.py

创建 replay/session.py

import asyncio, collections, contextlib, uuid
from enum import Enum

from .clock import ReplayClock, replay_batches, new_stats, summarize

class State(str, Enum):
    CREATED, RUNNING, PAUSED, COMPLETED, STOPPED = (
        "created", "running", "paused", "completed", "stopped"
    )

class ReplaySession:
    PRIORITY = {
        "paused", "resumed", "speed_changed",
        "replay_reset", "session_stopped"
    }

    def __init__(self, tape, speed=1.0, warmup_ms=120_000, maxsize=256):
        self.id = uuid.uuid4().hex[:12]
        self.tape = tape
        self.warmup_ms = warmup_ms
        self.maxsize = maxsize

        self.state = State.CREATED
        self.cursor = 0
        self.clock = ReplayClock(tape.ts[0], speed)
        self.stats = new_stats(speed)

        self._q = collections.deque()
        self._wake = asyncio.Event()
        self._task = None
        self._epoch = 0
        self._lock = asyncio.Lock()

    def info(self):
        lo, hi = self.tape.span()

        return {
            "session_id": self.id,
            "symbol": self.tape.symbol,
            "state": self.state.value,
            "speed": self.clock.speed,
            "cursor": self.cursor,
            "total_events": len(self.tape),
            "market_ts_ms": int(
                self.tape.ts[min(self.cursor, len(self.tape)-1)]
            ),
            "session_start_ms": lo,
            "session_end_ms": hi,
            "queued": len(self._q)
        }

    def _ctrl(self, kind, **kw):
        msg = {
            "type": kind,
            "session_id": self.id,
            "source": "replay",
            **kw
        }

        if kind in self.PRIORITY:
            self._q.appendleft(msg)
        else:
            self._q.append(msg)

        self._wake.set()

    async def _put(self, msg):
        while len(self._q) >= self.maxsize:
            self._wake.set()
            await asyncio.sleep(0)

        self._q.append(msg)
        self._wake.set()

    async def _kill(self):
        t, self._task = self._task, None

        if t and not t.done():
            t.cancel()

            with contextlib.suppress(
                asyncio.CancelledError,
                Exception
            ):
                await t

    async def start(self):
        self.clock.start()
        self.state = State.RUNNING
        self._task = asyncio.create_task(self._run())

        self._ctrl(
            "session_started",
            info=self.info()
        )

        return self.info()

    async def pause(self):
        if self.state is State.RUNNING:
            async with self._lock:
                self.clock.pause()
                self.stats["dropped"] += 1
                self.state = State.PAUSED

                self._ctrl(
                    "paused",
                    market_ts_ms=self.info()["market_ts_ms"]
                )

        return self.info()

    async def resume(self):
        if self.state is State.PAUSED:
            async with self._lock:
                self.clock.resume()
                self.state = State.RUNNING

                if self._task is None or self._task.done():
                    self._task = asyncio.create_task(self._run())

                self._ctrl(
                    "resumed",
                    market_ts_ms=self.info()["market_ts_ms"]
                )

        return self.info()

    async def set_speed(self, speed):
        async with self._lock:
            old = self.clock.speed
            self.clock.set_speed(speed)
            self.stats["speed"] = speed

            self._ctrl(
                "speed_changed",
                old_speed=old,
                new_speed=speed
            )

        return self.info()

    async def seek(self, target_ms):
        was = self.state
        await self._kill()

        async with self._lock:
            idx = max(
                0,
                min(
                    self.tape.index_at(target_ms),
                    len(self.tape)-1
                )
            )

            self._epoch += 1
            self.cursor = idx
            self.state = State.PAUSED
            self.clock.pause()

            warm = max(
                0,
                self.tape.index_at(
                    int(self.tape.ts[idx]) - self.warmup_ms
                )
            )

            purged = sum(
                1 for m in self._q
                if m.get("type") == "trade"
            )

            self._q = collections.deque(
                m for m in self._q
                if m.get("type") != "trade"
            )

            self._ctrl(
                "replay_reset",
                reason="seek",
                target_timestamp_ms=int(self.tape.ts[idx]),
                warmup_from_ms=int(self.tape.ts[warm]),
                warmup_events=idx-warm,
                purged_stale_events=purged,
                epoch=self._epoch
            )

        for k in range(warm, idx):
            await self._put({
                **self.tape[k].to_wire(),
                "warmup": True
            })

        self._ctrl(
            "warmup_complete",
            market_ts_ms=int(self.tape.ts[idx])
        )

        async with self._lock:
            self.clock.seek(float(self.tape.ts[idx]))

            if was is State.RUNNING:
                self.clock.start()
                self.state = State.RUNNING
                self._task = asyncio.create_task(self._run())

        return self.info()

    async def stop(self):
        self.state = State.STOPPED
        await self._kill()

        self._ctrl(
            "session_stopped",
            info=self.info(),
            timing=summarize(self.stats)
        )

        return self.info()

    async def _run(self):
        epoch = self._epoch

        async for i, j in replay_batches(
            self.tape,
            self.clock,
            self.cursor,
            self.stats
        ):
            if self._epoch != epoch or self.state is State.STOPPED:
                return

            for k in range(i, j):
                await self._put(self.tape[k].to_wire())
                self.cursor = k+1

        if self._epoch == epoch and self.cursor >= len(self.tape):
            self.state = State.COMPLETED

            self._ctrl(
                "session_completed",
                info=self.info(),
                timing=summarize(self.stats)
            )

    async def events(self):
        while True:
            if not self._q:
                self._wake.clear()
                await self._wake.wait()
                continue

            m = self._q.popleft()
            yield m

            if m.get("type") in (
                "session_completed",
                "session_stopped"
            ):
                return

会话状态的主要部分是 cursor,它指向 TradeTape 中的下一个位置。生产者使用时钟层的 replay_batches(),将每个到期的磁带位置转换为可通过网络传输的交易事件,并将其放入会话队列。

暂停会在不改变 cursor 的情况下冻结时钟。恢复会为时钟提供一个新的墙时间锚点,并从相同的历史位置继续。速度变化的工作方式类似:时钟首先在当前回放时间戳处锚定,然后从该点开始应用新的速度。

队列不仅包含交易。诸如 pausedresumedspeed_changedreplay_reset 之类的控制也会变成事件,这意味着下游消费者可以对回放状态的变化做出反应,而不必尝试从交易时间戳中推断这些变化。

seek() 是最复杂的控制。它会停止当前的生产者,使用 TradeTape.index_at() 找到请求的位置,移除过时的队列交易,并在播放继续之前准备一个预热窗口。我们将在有状态消费者就位后解释为什么需要这个预热。

目前,ReplaySession 没有独立的终端运行。在系统其余部分连接完成后,我们将通过实际的 API 和 WebSocket 流来演练这些控制。

使用 FastAPI 和 WebSocket 暴露回放

回放会话现在具备控制历史回放所需的一切,但它目前仅以 Python 对象的形式存在。为了让其他程序能够创建会话、控制会话并接收由此产生的交易流,我们将在其周围添加一个小的 API 层。

该层位于独立的 api/ 包中:

market-time-machine/
├── replay/
│   └── ...
└── api/
    ├── __init__.py
    ├── server.py
    └── run.py

我们将使用两条通信路径。REST 端点构成控制平面,而一个持久的 WebSocket 负责传输事件流。

Control plane

POST /sessions
POST /sessions/{id}/start
POST /sessions/{id}/pause
POST /sessions/{id}/resume
POST /sessions/{id}/speed
POST /sessions/{id}/seek
POST /sessions/{id}/stop


Event stream

WS /sessions/{id}/stream

诸如 pause 或 seek 之类的命令因此通过 HTTP 到达,而交易和回放控制事件则继续通过 WebSocket 流向消费者。

创建 api/server.py

创建 api/server.py

from fastapi import FastAPI, HTTPException, WebSocket, WebSocketDisconnect
from pydantic import BaseModel, Field

from replay import config
from replay.events import TradeTape
from replay.session import ReplaySession
from replay.clock import summarize

app = FastAPI(title="Market Time Machine")

SESSIONS = {}
ATTACHED = set()

class Create(BaseModel):
    symbol: str = "AAPL"
    date: str
    tag: str = "fullday"
    speed: float = Field(1.0, gt=0)
    warmup_ms: int = 120_000

class Speed(BaseModel):
    speed: float = Field(..., gt=0)

class Seek(BaseModel):
    target_timestamp_ms: int

def get(sid):
    if sid not in SESSIONS:
        raise HTTPException(404, f"no session {sid}")
    return SESSIONS[sid]

@app.post("/sessions")
async def create(b: Create):
    path = config.PROCESSED / f"{b.symbol}_{b.date}_{b.tag}.npz"

    if not path.exists():
        raise HTTPException(404, f"no tape {path.name}")

    s = ReplaySession(
        TradeTape.load(path),
        b.speed,
        b.warmup_ms
    )

    SESSIONS[s.id] = s
    return s.info()

@app.get("/sessions/{sid}")
async def info(sid: str):
    return get(sid).info()

@app.get("/sessions/{sid}/timing")
async def timing(sid: str):
    return summarize(get(sid).stats)

@app.post("/sessions/{sid}/start")
async def start(sid: str):
    return await get(sid).start()

@app.post("/sessions/{sid}/pause")
async def pause(sid: str):
    return await get(sid).pause()

@app.post("/sessions/{sid}/resume")
async def resume(sid: str):
    return await get(sid).resume()

@app.post("/sessions/{sid}/stop")
async def stop(sid: str):
    return await get(sid).stop()

@app.post("/sessions/{sid}/speed")
async def speed(sid: str, b: Speed):
    return await get(sid).set_speed(b.speed)

@app.post("/sessions/{sid}/seek")
async def seek(sid: str, b: Seek):
    return await get(sid).seek(b.target_timestamp_ms)

@app.websocket("/sessions/{sid}/stream")
async def stream(ws: WebSocket, sid: str):
    await ws.accept()

    if sid not in SESSIONS:
        return await ws.close(4004, "unknown session")

    if sid in ATTACHED:
        return await ws.close(4009, "consumer already attached")

    ATTACHED.add(sid)

    try:
        await ws.send_json({
            "type": "attached",
            "session_id": sid
        })

        async for msg in SESSIONS[sid].events():
            await ws.send_json(msg)

    except (WebSocketDisconnect, Exception):
        pass

    finally:
        ATTACHED.discard(sid)

创建会话会加载已处理的 .npz 磁带,并将其封装在 ReplaySession 中。此时,API 不再需要调用 EODHD 或重新读取原始 JSONL 响应。回放完全基于归一化的磁带进行。

REST 处理器故意保持精简。/pause 例如,不包含任何自身的暂停逻辑:

@app.post("/sessions/{sid}/pause")
async def pause(sid: str):
    return await get(sid).pause()

它只是将命令传递给 ReplaySession。同样的模式也适用于恢复、速度调整、查找和停止。这样可以把回放行为保持在 replay/ 目录内,而不是与 FastAPI 耦合。

WebSocket 端点负责处理另一方向。消费者连接后,服务器会将 session.events() 产生的所有数据转发过去:

async for msg in SESSIONS[sid].events():
    await ws.send_json(msg)

这可能是一笔正常的交易:

{
  "type": "trade",
  "symbol": "AAPL",
  "timestamp_ms": 1784122200009,
  "price": 317.46,
  "size": 3,
  "sequence": 61530328,
  "source": "replay"
}

或者是一条重放控制消息:

{
  "type": "paused",
  "market_ts_ms": 1784122200009
}

稍后,Seeking(跳转)还会引入另一个重要的控制事件:

{
  "type": "replay_reset",
  "reason": "seek",
  "target_timestamp_ms": 1784136600030
}

服务器每个回放会话仅支持一个 WebSocket 消费者。当前队列采用 FIFO 交接方式,而非广播机制,若在同一会话上挂载多个消费者,它们将共享事件流,而不是各自获得完整的事件流。

创建 api/run.py

第二个 API 文件仅用于启动 FastAPI 应用。

创建 api/run.py

import argparse
import uvicorn

from api.server import app

if __name__ == "__main__":
    p = argparse.ArgumentParser()
    p.add_argument("--port", type=int, default=8765)
    a = p.parse_args()

    uvicorn.run(
        app,
        host="127.0.0.1",
        port=a.port,
        log_level="warning"
    )

在项目根目录下启动服务:

python -m api.run --port 8765

回放引擎现在具备外部控制接口和 WebSocket 事件流。流的另一端是一个程序——它仅根据收到的事件来构建市场状态的消费者。

构建有状态的 WebSocket 消费者

回放服务现在可以流式传输历史交易,但我们仍然需要 WebSocket 另一端的某个组件,使其表现得像真实的下游应用。

该消费者不应直接读取历史磁带或调用 EODHD;它对市场的完整认识应仅来源于回放流中传来的消息。

我们将把它放在一个独立的包中:

market-time-machine/
├── replay/
│   └── ...
├── api/
│   └── ...
└── consumer/
    ├── __init__.py
    └── consumer.py

本教程中,消费者将维护:

该最终状态并不旨在作为交易策略。我们只需要一个真正具备状态的东西,以便以后验证回放控制(尤其是查找)不会让消费者保留过时的市场历史。

创建 consumer/consumer.py

创建 consumer/consumer.py

import argparse, asyncio, collections, json
import websockets

class VWAP:
    def __init__(self, window_ms):
        self.w = window_ms
        self.buf = collections.deque()
        self.pv = 0.0
        self.vol = 0.0

    def add(self, ts, px, sz):
        self.buf.append((ts, px, sz))
        self.pv += px * sz
        self.vol += sz

        cut = ts - self.w

        while self.buf and self.buf[0][0] < cut:
            _, p, s = self.buf.popleft()
            self.pv -= p * s
            self.vol -= s

        if self.vol <= 0:
            self.pv = self.vol = 0.0

    @property
    def value(self):
        return self.pv / self.vol if self.vol > 0 else None


class State:
    def __init__(self, short_ms=30_000, long_ms=120_000):
        self.short = VWAP(short_ms)
        self.long = VWAP(long_ms)

        self.last_trade = None
        self.last_sale = None
        self.signal = None

        self.n = 0
        self.vol = 0
        self.odd = 0
        self.zero = 0
        self.warming = False

    def apply(self, m):
        ts = m["timestamp_ms"]
        px = m["price"]
        sz = m["size"]
        meta = m["metadata"]

        self.short.add(ts, px, sz)
        self.long.add(ts, px, sz)

        self.last_trade = px

        if meta["last_sale_eligible"]:
            self.last_sale = px

        self.n += 1
        self.vol += sz
        self.odd += meta["odd_lot"]
        self.zero += meta["zero_size"]

        s = self.short.value
        l = self.long.value

        if s is not None and l is not None:
            self.signal = (
                "SHORT_ABOVE"
                if s > l
                else "SHORT_BELOW"
            )

    def line(self):
        f = lambda v: "--" if v is None else f"{v:.4f}"

        return (
            f"n={self.n:>7,} "
            f"vol={self.vol:>9,} "
            f"trade={f(self.last_trade):>9} "
            f"sale={f(self.last_sale):>9} "
            f"vwap30s={f(self.short.value):>9} "
            f"vwap2m={f(self.long.value):>9} "
            f"sig={self.signal or '--':<11} "
            f"odd={100*self.odd/max(1,self.n):4.1f}% "
            f"zero={100*self.zero/max(1,self.n):4.1f}%"
        )


async def run(url, every=3000):
    st = State()

    async with websockets.connect(
        url,
        max_size=None
    ) as ws:
        print("[consumer] connected", flush=True)

        async for raw in ws:
            m = json.loads(raw)
            t = m["type"]

            if t == "trade":
                st.apply(m)

                if not st.warming and st.n % every == 0:
                    print(
                        f"[consumer] {st.line()}",
                        flush=True
                    )

            elif t == "replay_reset":
                print(
                    f"[consumer] RESET -> "
                    f"{m['target_timestamp_ms']} "
                    f"({m['warmup_events']} warmup, "
                    f"{m['purged_stale_events']} purged)",
                    flush=True
                )

                st = State()
                st.warming = True

            elif t == "warmup_complete":
                st.warming = False

                print(
                    f"[consumer] WARM DONE {st.line()}",
                    flush=True
                )

            elif t in (
                "session_completed",
                "session_stopped"
            ):
                print(
                    f"[consumer] {t.upper()} "
                    f"{st.line()}",
                    flush=True
                )
                break

            else:
                print(
                    f"[consumer] {t}",
                    flush=True
                )


if __name__ == "__main__":
    p = argparse.ArgumentParser()

    p.add_argument(
        "--url",
        required=True
    )

    p.add_argument(
        "--every",
        type=int,
        default=3000
    )

    a = p.parse_args()

    asyncio.run(
        run(a.url, a.every)
    )

滚动 VWAP 窗口基于市场时间戳,而不是交易笔数。每笔新交易会同时进入两个窗口,随着回放时间推进,超过 30 秒或两分钟的观测值会被移除。

因此,消费者的状态会逐步演变:

关键在于,这些状态并非来源于原始的 TradeTape。消费者仅知晓已通过 WebSocket 传递的事件。

在回放时间向前推进时,这种做法运作良好。然而,进行定位(seek)时情况会变得复杂,因为若在不重置消费者的情况下移动回放光标,消费者将继续携带来自交易日错误时间点的状态。

使定位状态安全

定位不仅仅是移动回放光标那么简单。如果消费者已经在交易日的某一点构建了滚动状态,那么在不重置该状态的情况下跳转到其他位置,就会导致两段不同的市场历史混杂在一起。

假设消费者已经推进到 14:00。其两分钟 VWAP 仍保留着大约从 13:58 开始的交易。如果我们仅仅把回放光标移回到 13:30 并继续输出交易,那些后续观测值仍会驻留在内存中:

因此,回放在恢复正常播放之前,需要先重置下游状态,并围绕新的时间戳重建该状态。

重置并预热消费者

我们为 ReplaySession 添加的 seek() 方法已经能够处理这一流程。关键步骤首先是停止当前的生产者,并在磁带上定位请求的位置:

was = self.state
await self._kill()

async with self._lock:
    idx = max(
        0,
        min(
            self.tape.index_at(target_ms),
            len(self.tape)-1
        )
    )

    self._epoch += 1
    self.cursor = idx
    self.state = State.PAUSED
    self.clock.pause()

接下来,它会计算出目标时间前两分钟的预热点:

warm = max(0, self.tape.index_at(int(self.tape.ts[idx]) - self.warmup_ms))

我们使用两分钟,因为这与消费者维护的最长滚动窗口相匹配。重新播放该区间即可在新位置重建 30 秒和两分钟的 VWAP。

在发送这些预热交易之前,会话队列中等待的普通交易消息会被移除:

purged = sum(1 for m in self._q if m.get("type") == "trade")
self._q = collections.deque(m for m in self._q if m.get("type") != "trade")

随后,会话会发送一个显式的 replay_reset 事件:

self._ctrl(
    "replay_reset",
    reason="seek",
    target_timestamp_ms=int(self.tape.ts[idx]),
    warmup_from_ms=int(self.tape.ts[warm]),
    warmup_events=idx-warm,
    purged_stale_events=purged,
    epoch=self._epoch
)

消费者的应对方式是丢弃当前状态:

elif t == "replay_reset":
    st = State()
    st.warming = True

现在,会话就可以发送紧邻目标之前的历史交易了:

for k in range(warm, idx):
    await self._put({
        **self.tape[k].to_wire(),
        "warmup": True
    })

self._ctrl(
    "warmup_complete",
    market_ts_ms=int(self.tape.ts[idx])
)

这些交易通过与正常回放事件完全相同的 State.apply() 逻辑,但在 warmingTrue 时,消费者会抑制其常规输出。

因此,完整的 seek 流程如下:

检查重建后的状态

在完整会话运行中,我们暂停了回放并定位到 13:30。在此请求时间戳之后或恰好在此时间戳的第一个实际事件是 1784136600030

消费者收到:

旧的消费者状态已不存在,且 3,456 条历史交易已经在会话中的新点周围重建了两个滚动 VWAP 窗口。现在可以恢复正常的定时回放,而不会在 seek 边界处携带市场状态。

重放完整的 AAPL 交易日

所有部分现在已连接。全天的带子可以通过 FastAPI 进行控制,而独立的消费者仅能看到通过 WebSocket 到达的交易和控制事件。

在第一个终端中启动回放服务:

python -m api.run --port 8765

在最终运行中,我们将从 10x 开始,暂停市场,切换到 50x,恢复,再次暂停,定位到 13:30,重建消费者状态,最后以 400x 向收盘方向运行。

运行完整回放

将以下内容保存为项目根目录下的临时文件 demo.py。此脚本仅为演示的驱动程序。回放引擎和消费者仍位于我们已经构建的软件包中。

import asyncio, os, subprocess, sys
import httpx

BASE = "http://127.0.0.1:8765"
ROOT = os.getcwd()
SEEK_1330_MS = 1784136600000

async def demo():
    async with httpx.AsyncClient(base_url=BASE, timeout=120) as c:
        r = await c.post("/sessions", json={
            "symbol": "AAPL",
            "date": "2026-07-15",
            "tag": "fullday",
            "speed": 10.0
        })

        sid = r.json()["session_id"]

        consumer = subprocess.Popen([
            sys.executable,
            "-u",
            "-m",
            "consumer.consumer",
            "--url",
            f"ws://127.0.0.1:8765/sessions/{sid}/stream",
            "--every",
            "25000"
        ], cwd=ROOT)

        await asyncio.sleep(1.5)

        controls = [
            ("START @10.0x", f"/sessions/{sid}/start", None, 4),
            ("PAUSE", f"/sessions/{sid}/pause", None, 1.5),
            (
                "SPEED 50x while paused",
                f"/sessions/{sid}/speed",
                {"speed": 50.0},
                0.3
            ),
            ("RESUME", f"/sessions/{sid}/resume", None, 3),
            ("PAUSE", f"/sessions/{sid}/pause", None, 1),
            (
                "SEEK 13:30 while paused",
                f"/sessions/{sid}/seek",
                {"target_timestamp_ms": SEEK_1330_MS},
                3
            ),
            (
                "RESUME after seek",
                f"/sessions/{sid}/resume",
                None,
                3
            ),
            (
                "SPEED 400.0x to the close",
                f"/sessions/{sid}/speed",
                {"speed": 400.0},
                2
            )
        ]

        for label, path, payload, wait in controls:
            print(f"\n--- {label} ---")

            if payload is None:
                await c.post(path)
            else:
                await c.post(path, json=payload)

            await asyncio.sleep(wait)

        for _ in range(600):
            await asyncio.sleep(1)

            state = (
                await c.get(f"/sessions/{sid}")
            ).json()

            if state["state"] in ("completed", "stopped"):
                break

        print(
            f"\nfinal: {state['state']} "
            f"{state['cursor']:,}/{state['total_events']:,}"
        )

        print(
            "timing:",
            (
                await c.get(f"/sessions/{sid}/timing")
            ).json()
        )

        if consumer.poll() is None:
            consumer.terminate()

asyncio.run(demo())

另开一个终端运行:

python demo.py

消费者以独立进程启动,并在播放开始前连接到 WebSocket。

实际运行的开始如下:

这张技术文章配图显示了一个C语言程序的代码片段,用于控制一个消费者和一个生产者之间的数据传输

因此,会话可以在此被停止、以不同速度重新锚定,并在不重新启动回放的情况下恢复。

下一条命令直接跳转到 13:30:

这是在完整系统中发生的上一节所述的状态安全查找。消费者会丢弃其旧状态,处理 3,456 条预热事件,随后才从新的市场时间戳继续。

然后我们可以加速会话的剩余部分:

消费者会持续根据传入的事件更新其状态,直到会话到达磁带末尾:

这两个计数描述的是不同的内容。会话光标结束于 1,032,409/1,032,409,表明它已经到达全天磁带的末尾。消费者报告了 306,343 条事件,因为其状态在查找过程中被清除,随后从该点重建。查找还会跳过历史磁带的一部分,而不是实时流式传输每一条被跳过的交易。

在此次运行中,realized_speed 特意保持未设置,因为回放因暂停、速度变化和查找而多次重新锚定。单一的端到端速度比率无法有意义地描述一个在过程中故意改变时钟的会话。

这里重要的是,同一份历史磁带能够完整地经历控制序列,消费者在查找后重建其状态,并且回放持续进行直到会话结束。

测试回放引擎

全天运行表明系统能够完成完整的控制序列,但仅凭终端输出无法判断回放是否保持有序、是否遵守暂停边界,或者在查找后是否重建了正确的状态。

我们将使用之前创建的较小的 quiet15m 磁带来测试这些行为:

market-time-machine/
└── tests/
    ├── __init__.py
    ├── conftest.py
    └── test_replay.py

测试套件涵盖四个方面:事件顺序、回放时序、暂停/恢复行为以及寻址后的状态重建。

创建 tests/test_replay.py

创建 tests/test_replay.py:

import asyncio
import numpy as np
import pytest

from replay import config
from replay.events import TradeTape
from replay.session import ReplaySession
from replay.clock import ReplayClock, replay_batches, new_stats, summarize

TAPE = sorted(config.PROCESSED.glob("*_quiet15m.npz"))[0]

@pytest.fixture
def tape():
    return TradeTape.load(TAPE)

async def collect(sess, seconds):
    out = []

    async def drain():
        async for m in sess.events():
            out.append(m)

    t = asyncio.create_task(drain())
    await asyncio.sleep(seconds)
    return out, t


@pytest.mark.asyncio
async def test_ordering(tape):
    s = ReplaySession(tape, speed=500)
    out, t = await collect(s, 0.1)

    await s.start()
    await asyncio.sleep(2)
    await s.stop()
    t.cancel()

    trades = [
        m for m in out
        if m["type"] == "trade"
    ]

    assert len(trades) > 1000

    keys = [
        (m["timestamp_ms"], m["sequence"])
        for m in trades
    ]

    assert keys == sorted(keys)
    assert len(set(keys)) == len(keys)


@pytest.mark.asyncio
@pytest.mark.parametrize("speed", [10, 50, 100])
async def test_timing(tape, speed):
    end = int(
        np.searchsorted(
            tape.ts,
            tape.ts[0] + 60_000,
            side="right"
        )
    )

    clock = ReplayClock(tape.ts[0], speed).start()
    st = new_stats(speed)

    async for i, j in replay_batches(tape, clock, 0, st):
        if j >= end:
            break

    r = summarize(st)

    assert abs(r["speed_error_pct"]) < 5
    assert r["lateness_p95_ms"] < 50


@pytest.mark.asyncio
async def test_pause_resume(tape):
    s = ReplaySession(tape, speed=100)
    out, t = await collect(s, 0.05)

    await s.start()
    await asyncio.sleep(1)

    await s.pause()

    n = len([
        m for m in out
        if m["type"] == "trade"
    ])

    await asyncio.sleep(1)

    assert len([
        m for m in out
        if m["type"] == "trade"
    ]) == n

    await s.resume()
    await asyncio.sleep(1)

    await s.stop()
    t.cancel()

    seqs = [
        m["sequence"]
        for m in out
        if m["type"] == "trade"
    ]

    assert seqs == sorted(seqs)
    assert len(set(seqs)) == len(seqs)


@pytest.mark.asyncio
async def test_pause_seek_resume(tape):
    s = ReplaySession(
        tape,
        speed=200,
        warmup_ms=120_000
    )

    out, t = await collect(s, 0.05)

    await s.start()
    await asyncio.sleep(0.5)
    await s.pause()

    target = int(tape.ts[0]) + 300_000
    await s.seek(target)

    assert s.info()["state"] == "paused"

    def past():
        return [
            m for m in out
            if m["type"] == "trade"
            and not m.get("warmup")
            and m["timestamp_ms"] >= target
        ]

    await asyncio.sleep(0.4)
    assert not past()

    await s.resume()
    await asyncio.sleep(1)

    got = past()

    await s.stop()
    t.cancel()

    assert got

    seqs = [m["sequence"] for m in got]

    assert seqs == sorted(seqs)
    assert len(set(seqs)) == len(seqs)


@pytest.mark.asyncio
async def test_seek_state_equivalence(tape):
    import sys

    sys.path.insert(0, str(config.ROOT))
    from consumer.consumer import State as ConsumerState

    s = ReplaySession(
        tape,
        speed=200,
        warmup_ms=120_000
    )

    live = ConsumerState()
    reset = None
    snap = None
    out = []

    async def drain():
        nonlocal live, reset, snap

        async for m in s.events():
            out.append(m)

            if m["type"] == "trade":
                live.apply(m)

            elif m["type"] == "replay_reset":
                reset = m
                live = ConsumerState()

            elif m["type"] == "warmup_complete":
                snap = (
                    live.n,
                    live.vol,
                    live.short.value,
                    live.long.value
                )

    t = asyncio.create_task(drain())

    await s.start()
    await asyncio.sleep(1)

    await s.seek(
        int(tape.ts[0]) + 600_000
    )

    for _ in range(100):
        if snap:
            break
        await asyncio.sleep(0.05)

    await s.stop()
    t.cancel()

    assert snap

    fresh = ConsumerState()

    lo = tape.index_at(
        reset["warmup_from_ms"]
    )

    hi = tape.index_at(
        reset["target_timestamp_ms"]
    )

    for k in range(lo, hi):
        fresh.apply(tape[k].to_wire())

    n, vol, short, long = snap

    assert n == fresh.n == reset["warmup_events"]
    assert vol == fresh.vol

    assert short == pytest.approx(
        fresh.short.value,
        rel=1e-12
    )

    assert long == pytest.approx(
        fresh.long.value,
        rel=1e-12
    )

    kinds = [m["type"] for m in out]

    seg = out[
        kinds.index("replay_reset") + 1:
        kinds.index("warmup_complete")
    ]

    assert not [
        m for m in seg
        if m["type"] == "trade"
        and not m.get("warmup")
    ]

test_ordering() 检查发出的交易是否始终按 (timestamp, sequence) 排序,并且同一事件不会被发出两次。

时序测试在 10x、50x、100x 三种速度下运行 60 秒的历史市场时间。它允许一定的误差,而不是期望事件循环像硬实时调度器那样表现:实际速度必须保持在目标的 5% 以内,而第 95 百分位延迟必须低于 50 毫秒。

test_pause_resume() 验证另一种情况。在 pause() 返回后,已收到的交易数量应在播放恢复前保持不变。恢复后,得到的序列仍须保持有序且无重复。

test_pause_seek_resume() 覆盖完整回放中使用的确切控制模式。会话先暂停,然后将磁带前进五分钟,在新位置保持暂停状态,仅在调用 resume() 后才开始释放正常的寻址后交易。

独立验证状态重建

最强的测试是 test_seek_state_equivalence()

当回放进行寻址时,消费者会先收到一个重置,随后是两分钟的预热事件。此测试不只是检查是否出现 warmup_complete 消息,而是构建一个全新的 ConsumerState,并直接从磁带喂入相同的历史间隔,以独立验证:

for k in range(lo, hi):
    fresh.apply(tape[k].to_wire())

随后,重放构建的状态与独立重建的状态必须在以下几点上保持一致:

event count
cumulative volume
30-second VWAP
2-minute VWAP

VWAP 值会以相对容忍度 1e-12 进行比较。测试还会确保在 replay_resetwarmup_complete 之间没有普通回放交易混入流中。

配置 pytest

异步测试使用 pytest-asyncio。创建 tests/conftest.py

import pytest

def pytest_configure(config):
    config.addinivalue_line(
        "markers",
        "asyncio"
    )

然后在项目根目录下添加 pytest.ini 文件:

[pytest]
asyncio_mode = auto

运行完整的测试套件:

pytest tests/ -v

记录的运行产生了:

这些测试不仅检查回放是否最终到达磁带末尾,还验证历史顺序在回放中得以保留、加速后的时序是否在预期容忍范围内、控制是否保持事件序列,以及在查找后重建的状态是否与基于底层历史数据的独立重建相匹配。

结论

我最喜欢这个构建的地方在于,当我们再次为其添加时钟时,同样的历史数据集会呈现出完全不同的感觉。

我们从 EODHD 获得了一个完整的 AAPL 会话(EODHD),最终得到的结果可以缓慢前进、加速超前、在中途暂停、跳转到当天的另一个时间点,并在消费者仅对目前已接收到的内容做出反应时继续运行。

项目仍有很大的扩展空间。回放功能可以支持多个符号、更丰富的市场状态、多个下游消费者、持久的回放会话,甚至可以直接插入流中的策略和执行组件。当前版本故意将这些功能排除在外,但核心回放层已经具备,可作为后续开发的基础。

对我来说,这就是项目的有价值成果。EODHD 提供了历史事件,但回放层让另一段软件能够像体验交易日一样感受这些事件,而不是像一个已经知道当天如何结束的数据集。

——

🧑‍💻

zhirenhun

一个热爱技术的程序员,喜欢分享前沿AI知识和开发经验。

← 上一篇
如何自行基准测试LLM推理:值得信赖的数字设计标准
下一篇 →
GraphRAG 是推理问题,而非数据库问题

📌 相关推荐

GraphRAG 是推理问题,而非数据库问题
2026/8/30
如何自行基准测试LLM推理:值得信赖的数字设计标准
2026/8/30
如何从LLM中获取可靠的结构化数据
2026/8/28
← 返回文章列表