第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 实时特征计算框架

下面这张图,是我个人常用的实时特征计算架构:

数据源 交易所行情 数据清洗 去重/排序/异常检测 实时特征计算 滑动窗口/增量更新 特征存储 Redis / ClickHouse 策略引擎 定价/风控/下单 回测系统 历史数据重放 读取历史特征 离线批量清洗 实时特征计算架构图

这张图里,数据从交易所进来,经过清洗、实时计算,存入特征存储,然后被策略引擎和回测系统使用。注意看虚线部分——回测系统读取的是历史特征,而不是直接读原始数据,这样能保证训练和推理的一致性。

25.4 实战中的几个坑

最后,分享几个我在项目中踩过的坑,希望能帮你少走弯路:

  1. 时间同步问题:不同数据源的时间戳可能来自不同时钟。我建议统一用交易所的时间戳,别用本地时间。
  2. 特征膨胀:特征不是越多越好。我见过一个团队,搞了500多个特征,结果大部分是冗余的,反而增加了计算延迟。
  3. 冷启动问题:新策略上线时,滑动窗口里没有历史数据。我通常用「预热期」来解决,先跑一段时间再正式交易。
  4. 数据回放:回测时,要确保数据回放的速度和实时一致。我曾经用加速回放,结果特征计算逻辑和线上不一致,回测结果全是假的。

总结一句话:数据工程是做市商系统的地基。地基不稳,楼盖得再高也得塌。把数据清洗、特征存储、实时计算这三个环节做扎实了,你的策略才能跑得稳、跑得快。


交易系统化学习资料 微信Strategy888888