历史市场数据通常以完成的数据集形式到达。这对分析很方便,但与交易软件体验实时市场的方式截然不同。在生产环境中,事件会逐一到达,未来是未知的,每个决策仅依赖于迄今为止发生的事情。
在本教程中,我们将使用历史逐笔数据重建这种体验。我们将从 EODHD 获取完整的 AAPL 交易时段,将超过一百万笔交易归一化为确定性事件带,并通过可控的市场时钟按照原始时间顺序重播它们。
在此过程中,我们将添加可调的播放速度、暂停和恢复控制、搜索功能,以及一个通过 REST 暴露控制并通过 WebSocket 流式传输交易的 FastAPI 服务。
我们还将构建一个独立的消费者,仅根据其接收到的事件计算滚动 VWAP 和市场状态。完成后,我们将拥有一个完整的本地回放系统,能够将已经结束的交易日作为定时流喂送给事件驱动的软件,在搜索后正确重建下游状态,并通过自动化测试验证结果。
开始之前,请确保您已具备:
已安装 Python 3.10 或更高版本。
具有访问历史 tick-data 端点权限的 EODHD API 密钥。您可以在 EODHD 定价页面 创建开发者账户。
终端和代码编辑器。
基本的 Python 知识,包括函数、类、字典以及使用包。
对 HTTP 和 WebSocket 有基本了解。无需先前的 FastAPI 经验。
具有足够的本地磁盘空间来存储下载的原始 tick 数据和处理后的回放磁带。本教程中使用的完整 AAPL 会话包含超过一百万条交易记录。
本教程中的 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
整个构建的关键规则很简单:消费者只能知道回放流中已经到达的内容,永远不应该从历史磁带读取超前数据。正是这个约束,使得时序、暂停/恢复行为以及查找后的状态重构值得正确实现。
首先创建项目目录并安装我们将用于数据检索、回放时序、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
在 replay、api、consumer 和 tests 目录中创建空的 __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 的历史 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 获得的原始样子。但在回放引擎能够使用它们之前,它们需要转换为确定的内部事件序列。
加载器为我们提供了完整的会话,但回放引擎不应直接使用 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 的情况下冻结时钟。恢复会为时钟提供一个新的墙时间锚点,并从相同的历史位置继续。速度变化的工作方式类似:时钟首先在当前回放时间戳处锚定,然后从该点开始应用新的速度。
队列不仅包含交易。诸如 paused、resumed、speed_changed 和 replay_reset 之类的控制也会变成事件,这意味着下游消费者可以对回放状态的变化做出反应,而不必尝试从交易时间戳中推断这些变化。
seek() 是最复杂的控制。它会停止当前的生产者,使用 TradeTape.index_at() 找到请求的位置,移除过时的队列交易,并在播放继续之前准备一个预热窗口。我们将在有状态消费者就位后解释为什么需要这个预热。
目前,ReplaySession 没有独立的终端运行。在系统其余部分连接完成后,我们将通过实际的 API 和 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 另一端的某个组件,使其表现得像真实的下游应用。
该消费者不应直接读取历史磁带或调用 EODHD;它对市场的完整认识应仅来源于回放流中传来的消息。
我们将把它放在一个独立的包中:
market-time-machine/
├── replay/
│ └── ...
├── api/
│ └── ...
└── consumer/
├── __init__.py
└── consumer.py
本教程中,消费者将维护:
最新交易
最近的符合最后成交条件的交易
累计成交量
30秒 VWAP
2分钟 VWAP
零散股和零规模交易的百分比
一个简单的 SHORT_ABOVE / SHORT_BELOW 状态
该最终状态并不旨在作为交易策略。我们只需要一个真正具备状态的东西,以便以后验证回放控制(尤其是查找)不会让消费者保留过时的市场历史。
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() 逻辑,但在 warming 为 True 时,消费者会抑制其常规输出。
因此,完整的 seek 流程如下:
在完整会话运行中,我们暂停了回放并定位到 13:30。在此请求时间戳之后或恰好在此时间戳的第一个实际事件是 1784136600030。
消费者收到:
旧的消费者状态已不存在,且 3,456 条历史交易已经在会话中的新点周围重建了两个滚动 VWAP 窗口。现在可以恢复正常的定时回放,而不会在 seek 边界处携带市场状态。
所有部分现在已连接。全天的带子可以通过 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。
实际运行的开始如下:

因此,会话可以在此被停止、以不同速度重新锚定,并在不重新启动回放的情况下恢复。
下一条命令直接跳转到 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_reset 与 warmup_complete 之间没有普通回放交易混入流中。
异步测试使用 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 提供了历史事件,但回放层让另一段软件能够像体验交易日一样感受这些事件,而不是像一个已经知道当天如何结束的数据集。
——
一个热爱技术的程序员,喜欢分享前沿AI知识和开发经验。