第25章 做市商数据工程:数据清洗与预处理、特征存储、实时特征计算
做市商系统里,数据工程这块儿,说白了就是给交易模型「喂饭」的。饭好不好吃,干不干净,直接影响模型能不能跑起来。我见过太多团队,策略逻辑写得天花乱坠,结果数据源一塌糊涂,上线就崩。嗯,今天咱们就把这块硬骨头啃下来。
25.1 数据清洗与预处理:别让脏数据毁了你的策略
数据清洗,听着简单,做起来全是坑。我刚开始做市商那会儿,接过一个高频数据流,里面居然有重复的时间戳,还有负数的成交量。你想想看,这种数据喂给模型,它能给你什么好结果?
25.1.1 常见的数据脏问题
- 缺失值:比如某笔交易的卖一价突然为空。我建议用前向填充(ffill),别用均值填充,因为市场是序列相关的。
- 异常值:价格突然跳变几个数量级。这通常是数据源错误,直接剔除。
- 重复数据:同一个时间戳出现两笔相同数据。去重时注意保留最后一条。
- 时间戳错乱:数据到达顺序和实际发生顺序不一致。需要做排序和重采样。
核心原则:宁可丢掉一条可疑数据,也不要让一条脏数据进入特征计算。
25.1.2 预处理流水线示例
我个人习惯用 Python 搭一个轻量级流水线。下面这个例子,处理的是 Level 2 订单簿数据:
import pandas as pd
import numpy as np
def clean_orderbook(df):
# 1. 去重:按时间戳去重,保留最后一条
df = df.drop_duplicates(subset='timestamp', keep='last')
# 2. 排序:确保时间递增
df = df.sort_values('timestamp').reset_index(drop=True)
# 3. 处理缺失值:前向填充
df = df.ffill()
# 4. 异常值检测:价格超过3个标准差则剔除
price_cols = ['bid_price_1', 'ask_price_1']
for col in price_cols:
mean = df[col].mean()
std = df[col].std()
df = df[(np.abs(df[col] - mean) < 3 * std)]
# 5. 重采样到固定频率(比如100ms)
df.set_index('timestamp', inplace=True)
df = df.resample('100ms').last().dropna()
return df
避坑指南:我曾经在重采样时用了均值,结果把买卖价差算成了负数。记住,价格数据用 last(),成交量用 sum()。
25.2 特征存储:让数据「随取随用」
数据洗干净了,接下来就是存起来。特征存储(Feature Store)这个概念,说白了就是给所有特征一个统一的「家」。做市商系统里,特征数量动辄几百个,没有统一管理,后期维护就是噩梦。
25.2.1 为什么需要特征存储?
- 避免重复计算:同一个特征(比如「过去1分钟成交量」),多个策略都要用,算一次就够了。
- 保证一致性:离线训练和在线推理用同一套特征逻辑,不会出现「训练时用A,上线时用B」的尴尬。
- 版本管理:特征会迭代,旧版本需要保留用于回测。
25.2.2 存储架构设计
我推荐用分层存储:
| 层级 | 存储介质 | 用途 | 示例 |
|---|---|---|---|
| 热存储 | Redis / Memcached | 实时特征,毫秒级访问 | 当前买卖价差 |
| 温存储 | ClickHouse / InfluxDB | 分钟级特征,用于策略计算 | 过去5分钟波动率 |
| 冷存储 | Parquet / HDFS | 历史特征,用于回测和训练 | 过去30天所有特征 |
我的经验:热存储里只放当前时刻需要的特征,别把历史数据也塞进去。Redis 不是用来做时序数据库的。
25.3 实时特征计算:和时间赛跑
做市商的核心竞争力,就是比别人快。实时特征计算,要求你在数据到达后的几毫秒内,算出有用的信号。这里我分享几个实战技巧。
25.3.1 滑动窗口计算
最常见的实时特征就是滑动窗口统计量,比如「过去N笔交易的成交量加权平均价」。用纯 Python 算会很慢,我建议用 deque 或者 numpy 的滚动窗口:
from collections import deque
import numpy as np
class SlidingWindow:
def __init__(self, window_size):
self.window = deque(maxlen=window_size)
self.sum = 0.0
def update(self, value):
if len(self.window) == self.window.maxlen:
self.sum -= self.window[0]
self.window.append(value)
self.sum += value
return self.sum / len(self.window)
# 使用示例
vwap_calculator = SlidingWindow(100)
for price in realtime_price_stream:
vwap = vwap_calculator.update(price)
# 这里 vwap 就是实时计算的滑动平均
为什么用 deque? 因为它的 pop 和 append 都是 O(1) 复杂度。我曾经用 list 实现,数据量一上来,延迟直接飙到 50ms,换了 deque 后降到 1μs 以下。
25.3.2 增量更新 vs 全量重算
实时计算里,能增量更新就别全量重算。举个例子,计算「过去1分钟成交量」:
- 全量重算:每次请求都去查过去1分钟的所有数据,再求和。数据量大时,延迟不可控。
- 增量更新:维护一个计数器,新数据来了就加,旧数据过期了就减。延迟恒定在微秒级。
注意:增量更新需要处理好「过期数据」的移除。我见过一个系统,因为忘记移除过期数据,导致特征值越算越大,最后策略直接爆仓。
25.3.3 实时特征计算框架
下面这张图,是我个人常用的实时特征计算架构:
这张图里,数据从交易所进来,经过清洗、实时计算,存入特征存储,然后被策略引擎和回测系统使用。注意看虚线部分——回测系统读取的是历史特征,而不是直接读原始数据,这样能保证训练和推理的一致性。
25.4 实战中的几个坑
最后,分享几个我在项目中踩过的坑,希望能帮你少走弯路:
- 时间同步问题:不同数据源的时间戳可能来自不同时钟。我建议统一用交易所的时间戳,别用本地时间。
- 特征膨胀:特征不是越多越好。我见过一个团队,搞了500多个特征,结果大部分是冗余的,反而增加了计算延迟。
- 冷启动问题:新策略上线时,滑动窗口里没有历史数据。我通常用「预热期」来解决,先跑一段时间再正式交易。
- 数据回放:回测时,要确保数据回放的速度和实时一致。我曾经用加速回放,结果特征计算逻辑和线上不一致,回测结果全是假的。
总结一句话:数据工程是做市商系统的地基。地基不稳,楼盖得再高也得塌。把数据清洗、特征存储、实时计算这三个环节做扎实了,你的策略才能跑得稳、跑得快。