第14章:数据库存储方案:时序数据库选型、数据表结构设计、增量更新与数据校验
做基差交易,数据就是你的弹药库。没有靠谱的数据,再牛的策略也是空中楼阁。
我个人习惯把数据存储分成三层:原始数据层、清洗数据层、策略信号层。每一层对数据库的要求都不一样。今天咱们就聊聊,怎么把这套存储体系搭起来。
14.1 时序数据库选型:别用MySQL存K线
你可能想问:MySQL那么成熟,为什么不用?
我刚开始做量化时也踩过这个坑。用MySQL存了半年的1分钟K线,结果一张表就几千万行,查询速度慢得像蜗牛。更别提做回测时,要拉取几千只股票的历史数据,那简直是噩梦。
说白了,关系型数据库是为OLTP设计的,不是为时序数据设计的。时序数据的特点是:写多读少、按时间排序、数据量大。这时候,时序数据库(TSDB)才是正解。
我推荐几个常用的:
- InfluxDB:轻量级,单机部署方便,适合中小团队。我早期项目就用它,上手快,查询语法类似SQL。
- TimescaleDB:基于PostgreSQL的扩展,如果你团队熟悉SQL,这个过渡最平滑。我有个朋友在私募,他们就用这个,既保留了关系型数据库的灵活性,又有时序数据库的性能。
- ClickHouse:列式存储,查询速度极快,适合大规模数据分析。但部署和维护成本高一些,适合数据量在TB级别的场景。
我的建议:
如果数据量在100GB以内,用InfluxDB就够了。超过这个量,考虑TimescaleDB或ClickHouse。别一上来就上分布式,维护成本会让你崩溃。
14.2 数据表结构设计:基差数据的核心字段
设计表结构时,我习惯先想清楚:我要查什么?
基差交易中,最频繁的查询是:某时刻的基差是多少?某合约的基差历史走势?所以,时间戳和合约标识是必选字段。
下面是我常用的表结构,以InfluxDB为例:
-- 基差数据表
CREATE TABLE basis_data (
time TIMESTAMP, -- 时间戳,精确到毫秒
symbol STRING, -- 合约代码,如 'rb2401'
spot_price DOUBLE, -- 现货价格
futures_price DOUBLE, -- 期货价格
basis DOUBLE, -- 基差 = 现货 - 期货
basis_ratio DOUBLE, -- 基差率 = 基差 / 现货
volume INT, -- 成交量
open_interest INT, -- 持仓量
exchange STRING, -- 交易所
data_source STRING -- 数据来源
);
嗯,这里要注意:时间戳一定要用UTC。我见过有人用北京时间,结果夏令时切换时数据全乱了。统一用UTC,展示时再转换,这是铁律。
另外,我建议加一个data_source字段。为什么?因为不同数据源的数据可能有差异。我在项目中遇到过,同一个合约,Wind和Tushare的基差差了2个点。有了这个字段,你就能追溯问题。
14.3 增量更新:别每次都全量拉取
全量更新是新手最容易犯的错误。每次跑脚本,把几百万条数据全删了重新插入。这不仅慢,还容易把数据库搞崩。
正确的做法是增量更新。说白了,就是只拉取新增的数据。
具体怎么做?我分享一个实战方案:
- 记录上次更新时间:在数据库中维护一张
update_log表,记录每个合约的最后更新时间。 - 按时间戳查询:拉取数据时,只查询大于上次更新时间的数据。
- 批量插入:每500条或1000条一批,用批量插入语句,减少数据库连接开销。
代码示例(Python + InfluxDB):
from influxdb import InfluxDBClient
import pandas as pd
def incremental_update(symbol, last_update):
# 从数据源拉取增量数据
new_data = fetch_data_from_source(symbol, start_time=last_update)
if new_data.empty:
return
# 批量插入
client = InfluxDBClient(host='localhost', port=8086)
client.switch_database('basis_db')
json_body = []
for _, row in new_data.iterrows():
point = {
"measurement": "basis_data",
"tags": {"symbol": symbol},
"time": row['time'],
"fields": {
"spot_price": row['spot_price'],
"futures_price": row['futures_price'],
"basis": row['basis'],
"basis_ratio": row['basis_ratio'],
"volume": row['volume'],
"open_interest": row['open_interest']
}
}
json_body.append(point)
# 每500条写入一次
if len(json_body) >= 500:
client.write_points(json_body)
json_body = []
# 写入剩余数据
if json_body:
client.write_points(json_body)
# 更新日志
update_log(symbol, new_data['time'].max())
小技巧:
增量更新时,建议加一个dedup去重逻辑。因为网络波动可能导致重复数据。我通常用time + symbol作为唯一键,重复时覆盖更新。
14.4 数据校验:别让脏数据毁了你的策略
数据校验是最后一道防线。我曾经因为一个数据源的价格字段出现了空值,导致策略在回测时产生了巨额亏损。从那以后,我每次写入数据前都会做三件事:
- 空值检查:价格、基差等关键字段不能为空。
- 范围检查:价格不能为负数,基差率不能超过±20%(正常市场情况下)。
- 时间戳检查:时间戳不能是未来时间,也不能是重复时间。
代码示例:
def validate_data(df):
# 空值检查
if df['spot_price'].isnull().any() or df['futures_price'].isnull().any():
raise ValueError("价格字段存在空值")
# 范围检查
if (df['spot_price'] < 0).any() or (df['futures_price'] < 0).any():
raise ValueError("价格不能为负数")
if (df['basis_ratio'].abs() > 0.2).any():
raise ValueError("基差率异常,超过20%")
# 时间戳检查
if df['time'].max() > pd.Timestamp.now():
raise ValueError("存在未来时间戳")
return True
避坑指南:
我曾经遇到过一个问题:数据源在某个时间段内,所有合约的基差都变成了0。原因是数据源服务器宕机,返回了默认值。所以,我建议再加一个波动率检查:如果某个合约的基差在短时间内变化超过3个标准差,就触发告警,人工复核。
14.5 知识体系总览
下面这张图,是我对本章知识体系的总结。你可以把它当作一个检查清单:
这张图把整个存储方案串起来了。从选型到表结构,从增量更新到数据校验,每一步都环环相扣。你想想看,如果其中一环出了问题,后面的策略分析就全白费了。
好了,关于数据库存储方案,我就聊这么多。记住一句话:数据是量化交易的基石,别在存储上省钱省力。
公众号:蓝海资料掘金营,微信deep3321