构建一个“市场时间机器”:使用Python和WebSockets重新模拟交易过程
历史市场数据通常会以已完成的数据集的形式提供。这种格式便于进行分析,但它与交易软件体验实时市场的方式截然不同。在实际情况中,各种事件是依次发生的,未来的发展方向是不可预测的,因此每一个决策都只能基于迄今为止发生的事情来做出。 在本教程中,我们将利用历史成交数据来重现这种交易环境。我们会选取EODHD提供的完整AAPL交易日数据,将超过一百万笔交易转化为有序的事件序列,并通过一个可控制的市场时钟按照这些事件发生的原始时间顺序重新播放它们。 在这个过程中,我们会添加可调节的播放速度、暂停与继续播放的功能、搜索机制,还会开发一个FastAPI服务,该服务能够通过REST接口提供这些控制功能,同时通过
历史市场数据通常会以已完成的数据集的形式提供。这种格式便于进行分析,但它与交易软件体验实时市场的方式截然不同。在实际情况中,各种事件是依次发生的,未来的发展方向是不可预测的,因此每一个决策都只能基于迄今为止发生的事情来做出。
在本教程中,我们将利用历史成交数据来重现这种交易环境。我们会选取EODHD提供的完整AAPL交易日数据,将超过一百万笔交易转化为有序的事件序列,并通过一个可控制的市场时钟按照这些事件发生的原始时间顺序重新播放它们。
在这个过程中,我们会添加可调节的播放速度、暂停与继续播放的功能、搜索机制,还会开发一个FastAPI服务,该服务能够通过REST接口提供这些控制功能,同时通过WebSockets实时传输交易数据。
我们还会构建另一个专门用于处理数据的程序,它仅根据接收到的事件来计算滚动式VWAP值以及市场状态。最终,我们将拥有一个完整的本地回放系统,这个系统能够将已经结束的交易日的所有数据以有序流的形式传递给基于事件驱动机制的软件,并在用户进行搜索操作后正确地重建后续的市场状态,同时还会通过自动化测试来验证结果的准确性。
目录
已安装Python 3.10或更高版本。
拥有一个可以访问历史交易数据端点的EODHD API密钥。您可以通过EODHD定价页面创建开发账户。
需要一台终端和代码编辑器。
具备基本的Python知识,包括函数、类、字典以及如何使用包。
对HTTP和WebSockets有基本了解。无需具备FastAPI的使用经验。
拥有足够的本地磁盘空间来存储下载到的原始交易数据及处理后的回放数据。本教程中使用的完整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
整个系统设计的一个重要原则是:客户端只能知道通过回放流式传输系统接收到的数据,而绝不能提前读取历史数据。这一限制使得时间控制、暂停/恢复功能以及数据查找后的状态重建等功能能够得到正确实现。
设置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
在replay、api、consumer和tests目录下创建空的__init__.py文件,这样Python就会将这些目录视为独立的包:
replay/__init__.py
api/__init__.py
consumer/__init__.py
tests/__init__.py
我们将从EODHD获取历史交易数据,因此需要在项目根目录下创建一个.env文件,并将API密钥保存在其中:
下载到的交易数据文件体积会相当大,因此既不要将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的历史数据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未设置")
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,将数据请求的截止时间延长5秒。本教程中使用的示例数据包含了16:00之后的交易记录,因此这个缓冲时间能够确保这些数据也能被成功下载下来。
创建 replay/loader.py 文件
一次性发送大量请求并不足以安全地获取密集的行情数据。因为一天中交易活动会不断变化,如果某次请求所获取的数据量超过了配置的 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 rawexists() and man.exists() and not force:
m = json.loads(man.read_text())
print(f"已缓存 {raw.name} 文件:包含 {m['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"从 {cursor} 开始的区间包含超过 {config.MAX_LIMIT} 条记录,因此无法分页获取"
)
win = max(1, span // 2)
retries += 1
continue
if n > 0:
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 > 0 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:,} 条记录")
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("=== 全天交易会话数据 ===")
manifest, raw_path = fetch_session(
SYMBOL,
DATE,
tag="fullday"
)
print(
f" 时间窗口为 {manifest['window_from_utc']} 到 {manifest['window_to_utc']},
收集到了 {manifest['ticks']:,}条交易记录,共 {manifest['pages']}页数据,
经历了 {manifest['retries']}次重试尝试,
总耗时为 {manifest['elapsed_s']}秒,共进行了 {manifest['api_calls']}次API调用"
)
print(
f" 最早的交易时间戳为 {manifest['first_timestamp_ms']},
最晚的交易时间戳为 {manifest['last_timestamp_ms']}"
)
print(" 收集到的字段有:", manifest["fields"])
最终,在进行第二次数据采集时,系统重新使用了之前已经下载好的会话数据,从而得到了如下结果:
我们现在拥有了1,032,411条原始交易记录,这些数据涵盖了整个常规交易会话期间以及收盘后的附加时间窗口。其中145页数据被成功读取,而也经历了16次重试尝试,这些数据足以说明:对于如此密集的交易数据来说,使用固定的请求时间窗口显然是一个不合适的做法。
不过,这些记录仍然保持着它们从EODHD系统获取时的原始格式。在让回放引擎能够使用这些数据之前,还需要先对它们进行处理,以便将它们转换成一种统一的内部格式。
将交易数据转换为适合回放的格式
加载器虽然提供了完整的交易会话数据,但回放引擎直接使用EODHD提供的原始数据格式是无法正常工作的。因为EODHD返回的数据中,时间戳、价格、成交量、序列编号以及市场代码等字段都是以并行数组的形式存在的。
在回放这些数据之前,我们需要先确保这些数组中的数据是按照正确的顺序排列的,消除重复数据,并将最终结果转换成一种统一的内部格式。
相关的处理逻辑被放在了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"不完整的页面 {page['from']}: 数据结构不一致"
)
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.asarraycols["shares"], dtype=np.int64)
seq = np.asarray(cols["seq"], dtype=np.int64)
mkt = npассив(cols["mkt"], dtype=str)
sub = npассив(cols["sub_mkt"], dtype=str)
sl = npассив(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']:,} 条原始数据 -> {rep['kept']:,} 条有效数据 "
f"有 {rep['dupes']} 条重复数据,{rep['dropped']} 条无效数据"
)
print(
f"零交易量的记录占 {100*rep['zero_size']/n:.1f}% |
f"单笔交易量非整数的记录占 {100*odd/n:.1f}% |
f"符合最后成交条件的记录占 {100*elig/n:.1f}%"
)
print(
f"所有交易的序列号都是严格递增的:{rep['seq_strict']}"
)
return tape, rep
在第一步验证阶段,我们在构建任何交易记录之前就会进行检查。由于原始数据是以并行数组的形式提供的,因此页面上的每个字段都必须包含相同数量的观测值。否则,在将这些数据合并时,很可能会无意中将某笔交易的价格与另一笔交易的时间戳关联在一起。
之后,这些数组会被转换成NumPy格式,系统中会自动删除那些无效的数据记录,同时所有交易记录会按照(时间戳, 序列号)的顺序进行排序。时间戳确保了数据按时间顺序排列,而序列号则在多笔交易的时间戳相同的情况下提供了确定的排序规则。
接下来,系统会删除那些时间戳和序列号完全相同的重复记录。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"时间跨度为 {lo} 到 {hi} "
f"(相当于 {((hi-lo)/3_600_000:.2f} 个市场小时)"
)
print("前3条经过处理后的交易记录:")
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)
当播放速度为10倍时,500毫秒的时间间隔就会变为50毫秒;而当播放速度为100倍时,这个时间间隔则变为5毫秒。
问题在于asyncio.sleep()只能保证执行会在请求的延迟时间之后继续进行。如果每次等待的时间都略有偏差,并且下一次延迟时间是根据这次延迟后的时间来计算的,那么这些误差在长时间的回放过程中就会逐渐累积起来。
因此,我们应该将整个回放过程基于time.monotonic()来进行同步:
历史经过时间
÷
回放速度
+
墙钟起始时间
=
目标墙钟时间
这样一来,每个事件都会根据同一个基准点来安排执行顺序,而不是根据前一个事件结束的时间来确定。
创建replay/clock.py文件
将这个时钟模块添加到回放包中:
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(nppercentile(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()则执行相反的转换,用来确定某个历史时间戳在机器的单调时钟上应该对应什么时间。
通过暂停、恢复播放、调整播放速度或进行快进/倒带操作,我们可以重新设定这种时间映射关系,而无需修改原始的数据记录。
另一个重要的方面是数据分批处理。当回放速度很快时,如果为每一笔交易都安排一次延迟操作,将会产生额外的开销。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:,} 笔交易发生在AAPL quiet15m数据记录的前"
"30秒时间内\n"
)
async def measure():
print(
f"{'速度':>6} {'市场时间':>9} "
f"{'模拟时钟时间':>8} {'实际执行时间':>9} "
f"{'误差百分比':>7} {'50%延迟时间':>7} "
f"{'95%延迟时间':>7} {'最大延迟时间':>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['请求速度']:>6} "
f"{r['市场时间':>9} "
f"{r['模拟时钟时间':>8} "
f"{r['实际执行速度':>9} "
f"{r['误差百分比':>7} "
f"{r['50%延迟时间':>7} "
f"{r['95%延迟时间':>7} "
f"{r['最大延迟时间':>7}"]
)
await measure()
实际测试的结果如下:
在1倍速度下,处理30秒的历史市场数据需要30.001秒;在10倍速度下需要3.001秒;在50倍速度下需要0.6秒;而在100倍速度下只需要0.3秒。因此,实际执行速度与我们设定的目标值非常接近。
延迟数据可以告诉我们调度系统在处理各类交易时错过了多少时间。例如,在10倍速度下,中位延迟时间为0.36毫秒,95%的分位数延迟时间为1.12毫秒,而此次测试中最严重的延迟情况为11.75毫秒。
这些数值实际上用于测量“回放时钟”的运行状态。它们并不是端到端的WebSocket延迟测量结果,而且这种回放机制仍然是基于Python的事件循环进行最佳努力调度,并非采用交易所级别的计时系统。
通过回放会话添加播放控制功能
“回放时钟”能够知道交易应该在何时执行,但它并不知道当前的回放进度处于什么阶段,也不清楚是否应该开始进行回放操作。因此,我们需要另外一层机制来管理这些回放数据、跟踪当前的位置、处理事件队列,并协调诸如开始、暂停、恢复播放、调整播放速度、查找特定位置以及停止播放等控制操作。
这种逻辑应该被实现到replay/session.py文件中:
market-time-machine/
└── replay/
├── config.py
├── loader.py
├── events.py
├── clock.py
└── session.py
需要明确区分的是: “时钟”负责管理时间流逝,而“回放会话”则负责维护状态信息。
一个回放会话会经历一系列不同的状态变化:
创建replay/session.py
请创建replaysession.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
selfclock.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 = StateRUNNING
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流来测试这些控制功能。
使用FastAPI和WebSocket公开回放功能
目前,这个回放会话已经具备了控制历史数据回放所需的所有功能,但它仍然只以Python对象的形式存在。为了让其他程序能够创建这样的会话、对其进行控制,并接收回放产生的交易数据流,我们需要在它的外部添加一层API接口。
这一层API代码被放在一个独立的api/包中:
market-time-machine/
├── replay/
│ └── ...
└── api/
├── __init__.py
├── server.py
└── run.py
我们将使用两种通信方式:REST接口用于控制操作,而持久的WebSocket连接则用于传输事件数据流。
控制层
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
事件数据流
WS /sessions/{id}/stream
因此,像“暂停”或“定位”这样的指令是通过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"没有会话{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"没有对应的回放文件{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(id).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)
@appwebsocket("/sessions/{sid}/stream")
async def stream(ws: WebSocket, sid: str):
await ws.accept()
if sid not in SESSIONS:
return await ws.close(4004, "未知的会话")
if sid in ATTACHED:
return await ws.close(4009, "消费者已经连接")
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
}
将来,当用户进行定位操作时,还会产生另一种重要的控制消息:
{
"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在本教程中,该消费者会维护以下信息:
最新的交易记录
最近一次符合成交条件的交易记录
累计成交量
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 = 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] 已连接", 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] 重置操作完成 -> "
f"{m['target_timestamp_ms']} "
f"进行了{m['warmup_events']}次热身操作,"
f"清除了{m['purged_stale_events']}条过时数据",
flush=True
)
st = State()
st.warming = True
elif t == "warmup_complete":
st.warming = False
print(
f"[consumer] 热身操作已完成 {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秒或2分钟的数据,在重放过程进行过程中会被自动删除。
因此,消费者状态是逐步变化的:
重要的是,这些状态数据并非来源于最初的TradeTape数据;消费者只能了解到那些已经通过WebSocket传输过来的交易信息。
当重放过程继续进行时,这种处理方式能够确保数据的准确性。然而,在进行“回溯查询”操作时情况就会变得复杂了——因为如果在不重置消费者状态的情况下移动重放指针,那么消费者就会继续使用错误的时间点上的数据来进行计算。
确保回溯查询操作的准确性
进行“回溯查询”并不仅仅是移动重放指针这么简单。如果消费者在交易日的某个时刻已经生成了相应的状态数据,那么在没有重置这些数据的情况下直接跳到另一个时间点,就会导致两种不同的市场历史数据被混合在一起。
假设消费者已经查看了14:00这个时间点的交易数据,而它所计算的2分钟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))
我们选择2分钟作为这个“热身”间隔,是因为这个时间长度正好与消费者所维护的最长滚动窗口周期相匹配。重新播放这段时间内的交易数据,就足以在新的位置上重新计算出30秒和2分钟的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时,消费者会抑制自身的常规输出。
因此,整个寻址流程如下:
检查重建后的状态
在完整会话运行过程中,我们暂停了回放操作,并将时间点设置为13:30。在该时间点之后出现的第一个实际事件是1784136600030。
消费者接收到的数据如下:
之前的消费者状态已经被清除,而3,456条历史交易记录已经重新构建了会话中围绕新时间点产生的两个滚动式VWAP窗口。现在,正常的定时回放操作可以继续进行了,而且不会因为跨越寻址边界而导致市场状态出现混乱。
回放整个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:
consumerterminate()
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(setkeys)) == len(keys)
@pytest.mark.asyncio
pytest.mark.parametrize("speed", [10, 50, 100])
async def testtiming(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_errorpct"]) < 5
assert r["lateness_p95_ms"] < 50
@pytest.mark.asyncio
async def test_pauseresume(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()用于检查发出的交易是否仍然按照(时间戳, 序列号)的顺序排列,以及同一事件是否不会被重复发出。
这个计时测试会在历史市场时间的10倍、50倍和100倍速度下运行60秒。它允许存在一定的误差范围,并不要求事件循环的行为完全符合硬实时调度器的标准:实际执行速度必须保持在目标速度的5%以内,而95%置信区间的延迟时间也必须低于50毫秒。
test_pauseresume()测试的内容有所不同。当调用pause()函数后,在恢复播放之前,接收到的交易数量应该保持不变;恢复播放后,这些交易的顺序仍然需要是正确的,并且不能存在重复的交易。
testPauseSeekResume()则模拟了在完整回放过程中所使用的具体控制流程:先暂停回放,然后跳到5分钟前的位置,继续暂停在该位置,直到调用resume()之后,才会开始正常地发送后续的交易数据。
独立验证状态重建结果
其中最严格的测试是test_seek_state_equivalence()。
当回放系统进行定位操作时,消费者会首先收到一条重置信号,随后会接收到两分钟的分钟热身数据。这个测试并不只是简单地检查是否出现了warmup_complete消息,而是会重新创建一个全新的ConsumerState对象,然后直接从历史数据中读取相同的时间段内的交易数据来填充这个状态对象:
for k in range(lo, hi):
fresh.apply(tape[k].to_wire())
通过这种方式重新构建的状态和原始历史数据重建出的状态必须在以下方面保持一致:
交易数量
累计成交量
30秒移动平均成交量
2分钟移动平均成交量
这些移动平均成交量的数值会以1e-12的相对误差范围进行比较。此外,该测试还会确保在replay_reset和warmup_complete这两个时间点之间,没有正常的回放交易数据被误加入到数据流中。
配置pytest
这些异步测试使用了pytest-asyncio插件。需要创建一个名为tests/conftest.py的文件:
import pytest
def pytest_configure(config):
config.addinivalue_line(
"markers",
"asyncio"
)
接下来,在项目根目录下添加pytest.ini文件:
[pytest]
asyncio_mode = auto
最后运行完整的测试套件:
python pytest tests/ -v
这次测试产生的结果如下:
这些测试不仅检查回放过程是否能够顺利完成,还会验证历史数据的排序顺序在播放过程中是否得以保持,加速运行时的时间误差是否在允许的范围内,控制机制是否能够正确地维持事件序列的完整性,以及通过定位操作重建出的状态是否与从原始历史数据中独立计算得出的结果一致。结论
这个项目最让我印象深刻的地方在于:当我们再次为这些历史数据添加“时间线索”时,它们所呈现出的形态会变得截然不同。
我们最初使用的是来自EODHD的完整AAPL交易数据,而最终得到的系统则能够表现出如下特点:某些交易过程可能会进展缓慢,有时则会迅速发展;在某些时刻系统会暂停运行,然后跳转到当天的另一个时间点继续进行交易;而消费者的反应也只会针对那些已经发生过的交易事件。
这个项目还有很大的发展空间。未来的版本可以支持多个股票品种、更丰富的市场信息、多个下游应用系统,还可以实现持续性的回放功能;甚至还可以加入策略执行模块,让这些模块能够直接融入到数据流中。当前版本虽然有意省略了这些功能,但基础性的回放框架已经搭建完成,为后续的开发奠定了基础。
对我来说,这就是这个项目带来的真正价值。EODHD提供了历史交易数据,而回放功能则使得其他软件能够将这些数据视为一个真实的交易日来进行处理,而不是仅仅作为一个已经知道当天交易结果的数据集来使用。
相关文章
每位开发人员都应该了解的关于产品数据追踪的相关知识
产品数据记录了您的应用程序或网站内部实际发生的情况。它展示了用户的行为、系统的运行状态,以及业务的运营表现。 在本文中,您将了解什么是产品数据、哪些部分值得追踪、哪些部分可以忽略不计,同时也会明白:编写代码的开发者实际上承担着比他人认为的更大的责任。 目录 什么是产品数据? 为什么应该追踪产品数据? 还有谁会使用您所追踪的数据? 为什么应该尽早开始数据追踪? 为什么数据追踪永无止境? 在您的产品中应该追踪哪些内容? 有哪些内容是不应该被追踪的? 如何安全地处理用户数据? 总结 什么是产品数据? 产品数据指的是您的应用程序或网站内部实际发生的情况。它能够回答一些简单的问题:用户喜欢哪些功能?他们
阅读全文
程序化广告的运作原理
大多数关于 程序化广告 的教程都只停留在网页横幅这个层面。其实这很可惜,因为一旦你将这种概念应用到现实世界中,它会变得有趣得多。 如果你之前从未从事过广告行业的工作,也无需担心。你不需要任何广告行业的背景知识就能理解这些内容。只要你能阅读基本的Python代码,那就已经具备了所需的一切条件。 广告本身只是实现这一目标的一种手段而已。真正重要的技能是:如何将现实世界中那些杂乱无章的信息转化为软件能够处理的数据。 在这篇文章中,你将会了解什么是程序化广告,以及为什么它在网页上能够取得如此好的效果。你还会明白,为什么将这种技术应用到户外广告牌上实际上是一个数据建模的问题,并且你会通过编写简单的Pyt
阅读全文
使用OpenTelemetry实现Claude Code的可观测性
像 Claude Code 、 OpenAI Codex 、 Google Antigravity 以及 Cursor 这样的代理编码工具,在日常软件开发中已经变得无处不在。 随着代理系统的不断发展,开发者让这些系统完成的大部分工作都是通过逐个分配子任务来实现的。许多团队也在探索并使用共享的、多租户式的代理基础设施,这种架构的成本不会与某个特定的所有者挂钩。在这种情况下,可观测性就成为了监控基础设施成本的关键因素。 在本指南中,您将了解可观测性的工作原理,然后学习如何启用Claude Code内置的遥测功能,运行后端程序来收集数据,并读取该系统生成的各类指标、日志及追踪信息。这些内容将帮助您更
阅读全文
演示主题:如何通过一份数据同时完成从S3存储系统到GPU的计算任务——重新思考用于机器学习训练的数据加载方式
Onur Satici详细解释了如何利用Vortex这一由Linux基金会推出的开源列式文件格式来彻底改变高吞吐量数据加载的方式。他说明了如何通过层叠式的轻量级编码技术、基于文件结构的数据分段优化机制以及零拷贝内存传输技术,消除CPU与NVMe之间的性能瓶颈,从而让S3存储中的数据能够以高达60 Gbps的速度直接传输到GPU上,而无需在进行任何预处理操作。 作者:Onur Satici
阅读全文