市场数据回放系统交易模拟
本教程介绍了如何使用Python和WebSockets构建一个市场数据回放系统。通过从EODHD下载AAPL完整交易会话数据,将其标准化为事件磁带,并通过可控制的市场时钟进行回放。系统支持可调播放速度、暂停恢复、跳转等功能,并通过FastAPI提供REST API和WebSocket流式传输。还构建了有状态消费者计算滚动VWAP,最终实现完整的本地回放系统。
使用工具
为什么需要市场数据回放系统

在真实的交易环境中,市场数据是逐笔到达的:每一笔成交、每次报价都以事件流的形式推送给交易系统,系统只能根据已经发生的事件做出决策,而无法预知未来。在这种模式下,算法交易软件的运行逻辑与离线分析截然不同——历史数据集往往是完整可用的,而真实交易中“未知的未来”才是核心挑战。
为了让开发者能够更贴近真实市场体验地测试策略,我们需要构建一套市场数据回放系统。这种系统能够将历史 tick 数据按照原始时间顺序逐步“播放”,模拟真实市场中事件流的到达方式,从而支持交易模拟。
项目目标
- 从 EODHD 获取某只股票(如 AAPL)完整交易日的 tick 数据;
- 将超过一百万条交易记录标准化为确定性事件流;
- 通过可控时钟实现数据回放,支持播放速度调节、暂停/恢复、跳转等功能;
- 使用 FastAPI 提供 RESTful 接口,并通过 WebSocket 流式传输交易数据;
- 构建一个独立消费者,仅根据接收到的事件计算滚动 VWAP 和市场状态;
- 验证回放后状态的一致性,并编写自动化测试保障系统稳定。
搭建 Python 项目环境
首先确保本地安装了 Python 3.10 或更高版本,并创建一个虚拟环境:
python -m venv replay_env
source replay_env/bin/activate # Windows 用户使用 replay_env\Scripts\activate
pip install fastapi uvicorn websockets requests pytest
接下来,我们将逐步构建整个系统的各个模块。
获取完整交易日的历史数据
我们选择 EODHD 作为数据来源,因为它提供丰富的历史 tick 数据接口。通过其 API,可以下载某只股票某天内的所有交易记录。这些数据通常以 JSON 格式返回,包含时间戳、价格、成交量等字段。
在下载前,需要准备一个 EODHD 的开发者账号并获取 API Key。然后编写一个简单的数据加载器,用于请求并保存原始数据:
import requests
import json
def fetch_tick_data(symbol, date, api_key):
url = f"https://eodhd.com/api/tick-data/{symbol}?date={date}&api_key={api_key}"
response = requests.get(url)
data = response.json()
with open(f"{symbol}_{date}.json", "w") as f:
json.dump(data, f)
return data
定义配置文件与数据加载逻辑
为方便管理项目配置,我们创建 replay/config.py 文件,用于存放常量如 API Key、默认股票代码等:
# replay/config.py
API_KEY = "your_eodhd_api_key"
DEFAULT_SYMBOL = "AAPL"
DEFAULT_DATE = "2023-01-03"
然后在 replay/loader.py 中封装数据加载逻辑,使其可复用:
# replay/loader.py from .config import API_KEY, DEFAULT_SYMBOL, DEFAULT_DATE import json import os def load_or_fetch(symbol=DEFAULT_SYMBOL, date=DEFAULT_DATE): filename = f"{symbol}_{date}.json" if os.path.exists(filename): with open(filename, "r") as f: return json.load(f) else: return fetch_tick_data(symbol, date, API_KEY)
将 Tick 数据标准化为回放磁带
原始 tick 数据格式不一,且可能包含冗余信息。我们需要将其转化为一组统一结构的事件对象,形成所谓的“回放磁带”。每个事件应包含以下基本字段:
- timestamp: 事件发生时间;
- type: 事件类型(Trade / Quote);
- price: 价格;
- size: 成交量。
在 replay/events.py 中定义事件类并实现转换函数:
# replay/events.py
from dataclasses import dataclass
from typing import List, Dict
@dataclass
class TickEvent:
timestamp: str
event_type: str
price: float
size: int
def normalize_ticks(raw_data: List[Dict]) -> List[TickEvent]:
events = []
for record in raw_data:
events.append(TickEvent(
timestamp=record["timestamp"],
event_type=record["type"],
price=float(record["price"]),
size=int(record["size"])
))
return sorted(events, key=lambda x: x.timestamp)
构建可控的历史回放时钟
回放系统的核心是市场时钟,它决定了事件如何按照时间顺序被推送出去。我们设计一个可控时钟,允许用户调整播放速度、暂停、恢复甚至跳转到任意时间点。
在 replay/clock.py 中实现一个异步时钟类:
# replay/clock.py
import asyncio
from datetime import datetime, timedelta
from typing import List
from .events import TickEvent
class ReplayClock:
def __init__(self, events: List[TickEvent], speed: float = 1.0):
self.events = events
self.speed = speed
self.current_index = 0
self.paused = False
self.callbacks = []
async def start(self):
while self.current_index < len(self.events):
if self.paused:
await asyncio.sleep(0.1)
continue
event = self.events[self.current_index]
for cb in self.callbacks:
await cb(event)
self.current_index += 1
# 模拟时间间隔
await asyncio.sleep((1 / self.speed) * 0.001)
此外,我们还可以添加方法用于暂停、恢复、设置速度等操作。
添加播放控制与会话管理
为了更好地控制回放过程,我们引入 回放会话(Replay Session)概念。它封装了时钟的状态,并提供了更丰富的控制接口,如跳转到指定时间点、重置等。
在 replay/session.py 中实现会话类:
# replay/session.py
from .clock import ReplayClock
from .events import TickEvent
class ReplaySession:
def __init__(self, events: List[TickEvent]):
self.clock = ReplayClock(events)
self.listeners = []
def add_listener(self, listener):
self.listeners.append(listener)
self.clock.callbacks.append(listener.on_event)
async def run(self):
await self.clock.start()
def seek_to(self, timestamp: str):
target = datetime.fromisoformat(timestamp)
for i, event in enumerate(self.clock.events):
if datetime.fromisoformat(event.timestamp) >= target:
self.clock.current_index = i
break
使用 FastAPI 暴露回放服务
为了让其他系统能够远程访问回放数据,我们使用 FastAPI 构建一个 Web 服务。该服务不仅提供 REST 接口用于控制回放,还通过 WebSocket 实现实时数据推送。
在 api/server.py 中搭建服务:
# api/server.py
from fastapi import FastAPI, WebSocket
from replay.loader import load_or_fetch
from replay.session import ReplaySession
app = FastAPI()
session = None
@app.on_event("startup")
async def startup_event():
global session
raw_data = load_or_fetch()
events = normalize_ticks(raw_data)
session = ReplaySession(events)
@app.websocket("/ws")
async def websocket_endpoint(websocket: WebSocket):
await websocket.accept()
session.add_listener(WebSocketListener(websocket))
await session.run()
同时,在 api/run.py 中启动服务:
# api/run.py import uvicorn if __name__ == "__main__": uvicorn.run("api.server:app", host="0.0.0.0", port=8000) 构建状态感知型消费者
消费者负责接收回放事件,并据此计算市场状态,比如滚动 VWAP。它必须能够处理回放过程中可能出现的跳转行为,并在跳转后正确恢复状态。
在
consumer/consumer.py中实现消费者逻辑:# consumer/consumer.py class StatefulConsumer: def __init__(self): self.vwap_sum = 0.0 self.total_volume = 0 async def on_event(self, event: TickEvent): self.vwap_sum += event.price * event.size self.total_volume += event.size current_vwap = self.vwap_sum / self.total_volume if self.total_volume > 0 else 0 print(f"Current VWAP: {current_vwap}")当执行跳转操作时,消费者需要清空当前状态并重新播放从跳转点之前的所有事件,以保证状态一致性。
测试回放引擎
最后,我们编写单元测试验证整个系统的正确性。主要测试内容包括:
- 回放事件是否按时间顺序正确输出;
- 跳转后消费者状态是否正确恢复;
- WebSocket 是否正常推送数据。
在
tests/test_replay.py中编写测试用例:# tests/test_replay.py import pytest from replay.session import ReplaySession from replay.events import TickEvent @pytest.mark.asyncio async def test_replay_order(): events = [ TickEvent(timestamp="2023-01-03T09:30:00", event_type="Trade", price=150.0, size=100), TickEvent(timestamp="2023-01-03T09:31:00", event_type="Trade", price=151.0, size=200), ] session = ReplaySession(events) received = [] async def listener(event): received.append(event) session.add_listener(type('L', (), {'on_event': listener})()) await session.run() assert len(received) == 2同时配置
pytest.ini或pyproject.toml以启用异步测试支持。总结
通过本文所介绍的步骤,我们成功构建了一个完整的市场数据回放系统,它能够:
- 模拟真实市场中的事件驱动机制;
- 提供灵活的回放控制(速度调节、暂停、跳转);
- 通过 FastAPI + WebSocket 实现远程访问与实时推送;
- 构建具有状态感知能力的消费者,确保跳转后状态一致性;
- 通过自动化测试保障系统稳定可靠。
这套系统不仅适用于个人策略研发与回测,还可作为量化团队内部训练与评估工具使用,是构建高效交易模拟环境的理想基础。
相关推荐
利用Claude内置浏览器进行自动化网页任务执行
Anthropic为Claude桌面应用推出了内置浏览器功能,使AI能够直接在侧边栏加载网页、点击和输入。这使得用户能够自动化处理那些没有API接口的网站任务,如填写表单或抓取仪表盘数据,极大地提升了自动化办公效率。
无法确定自动化持续集成(CI)故障诊断机器人
本文通过48小时的实验,探讨了利用免费大模型自动化处理CI(持续集成)失败日志的方法。作者通过引入“确定性预过滤”机制,结合退出码和日志模式判断,有效解决了模型调用成本高、频率限制及幻觉问题,实现了高效的故障初步分诊。
不适用部署开源AI CEO进行业务自动化管理
该方法利用开源项目OpenExecutive构建虚拟CEO,通过AI决策引擎替代或辅助中小企业的管理决策。用户可通过n8n等工具将其集成到自动化工作流中,实现从日常决策到高层逻辑的自动化,旨在降低人力成本并提升决策效率,但需高度重视合规性与安全治理。
无法确定(取决于企业规模与自动化程度)本地PDF邮件合并自动化
该方法通过使用本地软件InOneShot,将Excel/CSV数据与PDF模板进行自动化合并,实现发票、证书等文档的批量生成。相比在线工具,该方案更注重数据隐私,无需上传敏感信息,适合小微企业和自由职业者通过自动化行政流程来节省大量时间。
不适用带有事实核查机制的自动化AI内容生成管线
本文通过一个因提示词中残留“$40”导致自动化发布失败的案例,深入探讨了构建高可靠性AI内容生成系统的核心:即建立“事实核查闸门(Grounding Gate)”。该方法强调通过严格的输入验证和输出溯源,防止AI幻觉,确保自动化内容生产的准确性。
未提及全自动AI博客内容流水线
本文作者通过构建全自动AI博客流水线进行实验,在1个月内生成了81篇文章。尽管实现了从选题到发布的完全自动化,但结果并不理想:Google仅索引了1篇文章,且搜索曝光量极低。作者强调通过API进行细粒度监控的重要性,并揭示了纯AI生成内容在SEO收录方面的巨大挑战。
未提及 (实验结果显示收入极低/接近于零)