市场数据回放系统交易模拟
本教程介绍了如何使用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 实现远程访问与实时推送;
- 构建具有状态感知能力的消费者,确保跳转后状态一致性;
- 通过自动化测试保障系统稳定可靠。
这套系统不仅适用于个人策略研发与回测,还可作为量化团队内部训练与评估工具使用,是构建高效交易模拟环境的理想基础。
相关推荐
自动化谈判跟进协议
本文介绍了一种针对自动化代理或自由职业者的谈判跟进协议。通过设定严格的触发条件(沉默24小时以上)和标准化的消息结构(价格锚定、明确范围、单一问题、拒绝预降价),旨在通过精准的跟进提高转化率,同时避免因过度跟进或过早让步而损害利润。
未提及利用无代码自动化优化自由职业工作流
本文分享了通过学习无代码自动化技术(如使用Zapier, Make, Airtable)来优化自由职业者工作流程的经验。通过将重复性的手动任务(如客户管理、发票处理)自动化,可以显著提升工作效率,打破业务增长的瓶颈。
未提及利用AI工具构建全自动化营销团队
本文介绍了如何利用五款低成本AI工具(ChatGPT, Midjourney, Buffer, Brevo, Canva)构建一个完整的营销团队,涵盖内容创作、视觉设计、社交媒体管理和邮件营销,旨在将原本每月数千美元的人力成本降低至不足100美元。
取决于具体业务规模 (文中强调的是节省成本,而非直接收入,但可用于降低运营成本)利用Seedeep监控Claude Code会话并优化成本
该内容介绍了一个名为Seedeep的开源工具,旨在为Claude Code提供可视化的监控界面。它能实时展示API调用延迟、Token消耗(区分缓存与新Token)、子代理运行状态及错误原因。通过该工具,开发者可以清晰识别Token浪费,优化上下文管理,从而显著降低使用Claude Code时的API账单成本。
不适用利用Claude插件实现自动化简化技术英语(STE)内容生成
该内容介绍了一种名为 SHOOK 的技术工具,通过为 Claude Code 开发自动化钩子(Hooks),强制 AI 遵循 ASD-STE100 简化技术英语标准。它通过规则注入、提示词提醒和 Lint 校验门禁,确保 AI 生成的内容始终符合专业技术文档的简洁性要求。这主要是一个提高技术写作效率的工具,而非直接的赚钱方法。
不适用WikiSkill AI智能体技能进化框架
Google Research推出的WikiSkill是一种通过持久化知识库提升AI智能体性能的框架。它通过“原始层-维基层-技能层”三层架构,让智能体能从过去的错误和成功中学习,将经验转化为可复用的“技能模块”,从而在不重新训练模型的情况下实现能力的持续进化。
不适用