29、订单流与高频数据清洗:Tick数据的清洗与对齐、订单簿数据的重建、订单流数据的存储与压缩
做订单流分析,最头疼的其实不是策略本身,而是数据。我见过太多人,策略逻辑写得漂漂亮亮,一上实盘就崩,最后发现是数据源出了问题。说白了,高频数据清洗是订单流交易的「地基」,地基不稳,楼盖得再高也得塌。
今天我们就来聊聊这个地基怎么打。我会从Tick数据清洗、订单簿重建,一直讲到数据存储与压缩。嗯,都是我在实战中踩过的坑,希望能帮你少走弯路。
一、Tick数据的清洗与对齐
Tick数据,就是交易所每笔成交的原始记录。它包含时间戳、价格、成交量、买卖方向等字段。但问题在于,不同交易所、不同数据源,给的Tick数据格式五花八门,而且经常有脏数据。
1.1 常见脏数据问题
- 时间戳错乱:比如某笔成交的时间戳比前一笔还早,或者时间戳精度不一致(有的到毫秒,有的到微秒)。
- 价格异常:比如出现0价格、负价格,或者价格超出当日涨跌停板范围。
- 成交量异常:比如成交量为0,或者成交量突然暴增(可能是数据拼接错误)。
- 重复数据:同一笔成交被推送了两次。
- 缺失数据:某些时间段的Tick数据完全丢失。
核心原则:清洗Tick数据,不是把异常数据删掉就完事了。你需要判断:这个异常是数据源的问题,还是市场真实的极端情况?比如某笔成交价格突然跳空,可能是数据错误,也可能是发生了闪电崩盘。
1.2 我的清洗流程
我个人习惯用以下步骤来处理Tick数据。你想想看,如果每一步都做到位,数据质量基本就有保障了。
- 时间戳标准化:把所有时间戳统一到同一个时区(通常是UTC),并统一精度(比如全部转为毫秒级时间戳)。
- 排序与去重:按时间戳排序,然后检查相邻两条数据的时间戳和成交ID。如果成交ID相同,说明是重复数据,保留第一条即可。
- 价格合理性检查:设定一个合理的价格范围(比如当日涨跌停价的±5%),超出范围的标记为异常,需要人工复核。
- 成交量检查:成交量为0的直接剔除。如果某笔成交量突然比前100笔的平均成交量高出10倍以上,我会标记为「可疑」,但不直接删除——因为可能是大单成交。
- 缺失数据插补:如果某段时间缺失数据,且缺失时间较短(比如几秒钟),可以用前一笔Tick的价格填充。如果缺失时间较长,建议直接丢弃该时间段。
避坑指南:我曾经遇到过一个数据源,它的Tick数据时间戳是本地时间,但没标注时区。我一开始没注意,直接用UTC处理,结果策略在非交易时段疯狂开仓。后来我花了整整两天才定位到这个问题。所以,拿到数据的第一件事,就是确认时间戳的时区和精度。
1.3 代码示例:Tick数据清洗
下面是一个简单的Python示例,演示了如何清洗Tick数据。注意,这只是个框架,实际生产中你需要根据数据源的具体格式来调整。
import pandas as pd
import numpy as np
def clean_tick_data(df):
"""
清洗Tick数据
df: DataFrame,包含列 ['timestamp', 'price', 'volume', 'side']
"""
# 1. 时间戳标准化:转为毫秒级时间戳
df['timestamp'] = pd.to_datetime(df['timestamp']).astype(np.int64) // 10**6
# 2. 排序
df = df.sort_values('timestamp').reset_index(drop=True)
# 3. 去重:基于时间戳和价格去重
df = df.drop_duplicates(subset=['timestamp', 'price', 'volume'])
# 4. 价格合理性检查
# 假设当日合理价格区间为 [100, 200]
df = df[(df['price'] >= 100) & (df['price'] <= 200)]
# 5. 成交量检查:剔除成交量为0的数据
df = df[df['volume'] > 0]
# 6. 缺失数据检查:如果时间间隔超过5秒,标记为缺失
df['time_diff'] = df['timestamp'].diff()
missing_mask = df['time_diff'] > 5000 # 5秒 = 5000毫秒
print(f"发现 {missing_mask.sum()} 个缺失时间段")
return df
二、订单簿数据的重建
有了干净的Tick数据,下一步就是重建订单簿。订单簿,说白了就是买卖双方的挂单队列。它由多个档位(Level)组成,每个档位包含价格和挂单量。
重建订单簿有两种方式:
- 基于快照+增量:交易所会定期推送订单簿快照(比如每100ms一次),同时推送增量更新(比如某价格挂单量变化)。你只需要在快照基础上,应用增量更新即可。
- 基于Tick数据反推:如果只有Tick数据,没有订单簿快照,那就只能通过Tick数据中的买卖方向,结合一些假设来反推订单簿。这种方法精度较低,但聊胜于无。
2.1 基于快照+增量的重建
这是最常用的方法。我个人建议,如果你能拿到订单簿快照,就尽量用这种方法。它准确、高效。
重建的核心逻辑是:
- 收到一个快照,清空当前订单簿,用快照数据填充。
- 收到增量更新,根据更新类型(新增、修改、删除)来调整对应价格档位的挂单量。
- 注意:增量更新可能乱序到达,所以需要根据时间戳或序列号来排序。
注意:有些交易所的增量更新是「全量替换」某个价格档位,有些是「增量调整」。一定要仔细阅读交易所的API文档。我见过有人把增量当成全量来用,结果订单簿越重建越离谱。
2.2 代码示例:订单簿重建
class OrderBook:
def __init__(self):
self.bids = {} # 买盘,key=价格,value=挂单量
self.asks = {} # 卖盘,key=价格,value=挂单量
def apply_snapshot(self, snapshot):
"""应用快照"""
self.bids.clear()
self.asks.clear()
for level in snapshot['bids']:
self.bids[level['price']] = level['volume']
for level in snapshot['asks']:
self.asks[level['price']] = level['volume']
def apply_update(self, update):
"""应用增量更新"""
side = self.bids if update['side'] == 'bid' else self.asks
price = update['price']
volume = update['volume']
if volume == 0:
# 删除该价格档位
side.pop(price, None)
else:
# 新增或修改
side[price] = volume
def get_top_n(self, n=5):
"""获取前N档买卖盘"""
sorted_bids = sorted(self.bids.items(), reverse=True)[:n]
sorted_asks = sorted(self.asks.items())[:n]
return sorted_bids, sorted_asks
三、订单流数据的存储与压缩
高频数据量非常大。一个交易所一天的Tick数据,可能就有几个GB。如果存储原始数据,一年下来就是TB级别。所以,存储和压缩是必须考虑的问题。
3.1 存储格式选择
| 格式 | 优点 | 缺点 | 适用场景 |
|---|---|---|---|
| CSV | 通用、易读 | 体积大、读写慢 | 小规模数据、调试 |
| Parquet | 列式存储、压缩率高、读写快 | 不直观、需要特定库 | 大规模数据、分析 |
| HDF5 | 支持复杂数据结构、压缩 | 单线程写入慢 | 科学计算、回测 |
| 数据库(如ClickHouse) | 支持实时查询、分布式 | 运维成本高 | 生产环境、实时分析 |
我个人比较推荐Parquet格式。它压缩率高,而且支持按列读取,回测时只需要读取需要的字段,速度很快。
3.2 数据压缩技巧
除了格式本身的压缩,我们还可以从数据层面做优化:
- 价格归一化:将价格转换为整数(比如乘以10000),然后存储为int32,比float64省一半空间。
- 时间戳差分:存储时间戳的差值,而不是完整时间戳。因为相邻Tick的时间差通常很小,用int16就能存下。
- 成交量归一化:如果成交量单位是手,可以存储为int32。如果单位是股,可能需要int64。
- 只存储必要字段:比如回测只需要价格和成交量,那就不要存储买卖方向、成交ID等字段。
避坑指南:我曾经为了省空间,把价格归一化后存为int32,结果没注意某个品种的价格超过了int32的范围(比如比特币价格几万美元,乘以10000后超过了21亿)。回测时数据全部溢出,策略表现一塌糊涂。所以,归一化之前一定要确认数据范围。
3.3 代码示例:数据压缩存储
import pandas as pd
import numpy as np
def compress_and_save(df, filepath):
"""
压缩并存储Tick数据
"""
# 价格归一化:乘以10000,转为int32
df['price_int'] = (df['price'] * 10000).astype(np.int32)
# 时间戳差分:计算与前一条的时间差
df['timestamp_diff'] = df['timestamp'].diff().fillna(0).astype(np.int16)
# 只保留必要字段
df_compressed = df[['timestamp_diff', 'price_int', 'volume']]
# 保存为Parquet
df_compressed.to_parquet(filepath, compression='snappy')
print(f"数据已保存到 {filepath},原始大小: {df.memory_usage().sum() / 1024**2:.2f} MB,压缩后: {df_compressed.memory_usage().sum() / 1024**2:.2f} MB")
四、知识体系总览
下面这张图,是我自己梳理的订单流数据清洗与存储的完整流程。你可以把它当作一个检查清单,每次处理数据时对照着来。
嗯,以上就是订单流数据清洗与存储的核心内容。从Tick数据清洗,到订单簿重建,再到数据压缩存储,每一步都有讲究。我个人觉得,数据清洗这部分花再多时间都值得——因为干净的数据,是策略盈利的前提。
最后提醒一句:不要迷信任何数据源。即使是交易所官方数据,也可能有bug。我见过某交易所的订单簿快照,在某个价格档位上突然多了一个0,导致挂单量暴增10倍。如果你不做校验,直接拿这个数据去跑策略,后果可想而知。
交易系统化学习资料 微信Strategy888888