主题切换
17.1 Tick 数据库设计
名词解释:Tick 数据
Tick 数据是金融市场中最细粒度的交易数据,记录每一笔成交(Trade)和订单簿每一次变动(Quote/Order Book Update)。与分钟级或日级别的 K 线数据不同,Tick 数据保留了完整的市场微观结构信息,是做市策略、高频交易和执行算法研究的必需品。Tick 数据的核心特点是高频、海量和时间敏感。
一、Tick 数据特点:高频、海量、时间敏感
1.1 数据量估算
以美国权益市场为例:
| 数据维度 | 日盘(6.5小时) | 备注 |
|---|---|---|
| 成交笔数 | 5000万~5亿笔/日 | 取决于市场活跃度 |
| 报价更新 | 5亿~50亿次/日 | SIP+直连交易所 |
| 原始数据量 | 50GB~500GB/日(压缩前) | 含所有字段 |
| 压缩后数据量 | 5GB~50GB/日 | 列式存储+压缩 |
全年约 252 个交易日,每只股票的 Tick 数据量相当可观。对于全市场历史数据,轻松达到 PB 级别。
1.2 Tick 数据的核心字段
此处有展示代码展开 ▼
python
# Tick 数据的标准化 Schema
TICK_TRADE_SCHEMA = {
'timestamp': 'datetime64[ns]', # 成交时间戳(纳秒精度)
'symbol': 'str', # 股票代码
'exchange': 'str', # 交易所代码
'price': 'float64', # 成交价
'size': 'int64', # 成交量
'trade_id': 'int64', # 成交编号
'trade_condition': 'str', # 成交条件(如 @, F, I 等)
'bid_price': 'float64', # 成交时的最优买价
'ask_price': 'float64', # 成交时的最优卖价
}
TICK_QUOTE_SCHEMA = {
'timestamp': 'datetime64[ns]', # 报价时间戳
'symbol': 'str',
'exchange': 'str',
'bid_price': 'float64', # 最优买价
'bid_size': 'int64', # 最优买价挂单量
'ask_price': 'float64', # 最优卖价
'ask_size': 'int64', # 最优卖价挂单量
'bid_exchange': 'str', # 最优买价来源交易所
'ask_exchange': 'str',
'quote_condition': 'str',
'nbbo_indicator': 'bool', # 是否为 NBBO
}
# 按 dtype 估算单条记录的字节数,进而推算全市场存储规模
DTYPE_BYTES = {
'datetime64[ns]': 8, 'float64': 8, 'int64': 8, 'bool': 1,
'str': 12, # 字典编码/LowCardinality 后的近似值
}
def schema_report(name, schema):
row_bytes = sum(DTYPE_BYTES[t] for t in schema.values())
print(f"=== {name}({len(schema)} 字段,{row_bytes} bytes/行)===")
for field, dtype in schema.items():
print(f" {field:<18}{dtype:<18}{DTYPE_BYTES[dtype]:>3} B")
print()
return row_bytes
trade_row = schema_report('TICK_TRADE_SCHEMA', TICK_TRADE_SCHEMA)
quote_row = schema_report('TICK_QUOTE_SCHEMA', TICK_QUOTE_SCHEMA)
# 存储量估算:A 股全市场约 5000 只标的
N_SYMBOLS, N_DAYS = 5000, 252
TRADES_PER_SYMBOL_DAY = 20_000 # 活跃标的日均成交笔数量级
QUOTES_PER_SYMBOL_DAY = 200_000 # 报价更新远多于成交
COMPRESS_RATIO = 15 # 列式存储 + delta/字典编码的典型压缩比
raw_trade = trade_row * TRADES_PER_SYMBOL_DAY * N_SYMBOLS * N_DAYS
raw_quote = quote_row * QUOTES_PER_SYMBOL_DAY * N_SYMBOLS * N_DAYS
raw_total = raw_trade + raw_quote
print(f"=== 全市场一年 Tick 存储估算({N_SYMBOLS} 标的 × {N_DAYS} 交易日)===")
print(f" 成交流水 原始: {raw_trade / 1e12:>7.2f} TB")
print(f" 报价快照 原始: {raw_quote / 1e12:>7.2f} TB")
print(f" 合计 原始: {raw_total / 1e12:>7.2f} TB")
print(f" 列式压缩后({COMPRESS_RATIO}x): {raw_total / COMPRESS_RATIO / 1e12:>7.2f} TB")
print(f"\n结论:报价流量约为成交流量的 {raw_quote / raw_trade:.0f} 倍,"
f"quote 表才是存储与查询成本的主要来源。")点击展开可浏览运行结果
=== TICK_TRADE_SCHEMA(9 字段,84 bytes/行)=== timestamp datetime64[ns] 8 B symbol str 12 B exchange str 12 B price float64 8 B size int64 8 B trade_id int64 8 B trade_condition str 12 B bid_price float64 8 B ask_price float64 8 B === TICK_QUOTE_SCHEMA(11 字段,101 bytes/行)=== timestamp datetime64[ns] 8 B symbol str 12 B exchange str 12 B bid_price float64 8 B bid_size int64 8 B ask_price float64 8 B ask_size int64 8 B bid_exchange str 12 B ask_exchange str 12 B quote_condition str 12 B nbbo_indicator bool 1 B === 全市场一年 Tick 存储估算(5000 标的 × 252 交易日)=== 成交流水 原始: 2.12 TB 报价快照 原始: 25.45 TB 合计 原始: 27.57 TB 列式压缩后(15x): 1.84 TB 结论:报价流量约为成交流量的 12 倍,quote 表才是存储与查询成本的主要来源。
二、列式存储(Parquet/ClickHouse)设计
2.1 列式存储的优势
对于 Tick 数据查询的典型模式(按时间范围过滤、按股票分组、聚合计算),列式存储相比行式存储有以下优势:
| 特性 | 行式存储 (PostgreSQL/MySQL) | 列式存储 (Parquet/ClickHouse) |
|---|---|---|
| 压缩比 | 2-5x | 10-50x(同类型数据压缩效率高) |
| 列选择查询 | 需要扫描整行 | 只读取需要的列 |
| 聚合查询 | 慢(需要全表扫描) | 快(向量化执行 + SIMD) |
| 时间范围过滤 | 如果有索引中等 | 快(分区裁剪 + min/max 统计) |
2.2 Parquet 文件的分区设计
此处有展示代码展开 ▼
python
import pyarrow as pa
import pyarrow.parquet as pq
import pandas as pd
import numpy as np
from pathlib import Path
from datetime import datetime, timedelta
def write_tick_to_parquet(tick_data: pd.DataFrame,
base_path: str,
partition_by: str = 'date') -> list:
"""
将 Tick 数据写入分区 Parquet 文件。
推荐分区方案:
/data/trades/{symbol}/{date}/{symbol}_{date}_{hour}.parquet
/data/quotes/{symbol}/{date}/{symbol}_{date}_{hour}.parquet
参数:
tick_data: Tick 数据 DataFrame(含 symbol, timestamp 列)
base_path: 存储根目录
partition_by: 分区方式 'date' 或 'symbol_date'
返回:
写入的文件路径列表
"""
written_files = []
# 添加日期列用于分区
tick_data['date'] = tick_data['timestamp'].dt.date
tick_data['hour'] = tick_data['timestamp'].dt.hour
# 按 (symbol, date) 分组
for (symbol, date), group in tick_data.groupby(['symbol', 'date']):
# 文件路径
date_str = date.strftime('%Y%m%d')
symbol_dir = Path(base_path) / symbol / date_str
symbol_dir.mkdir(parents=True, exist_ok=True)
# 按小时切分文件(控制文件大小)
for hour, hour_group in group.groupby('hour'):
filename = f"{symbol}_{date_str}_{hour:02d}.parquet"
filepath = symbol_dir / filename
# 写 Parquet,使用 Snappy 压缩
table = pa.Table.from_pandas(
hour_group.drop(columns=['date', 'hour'])
)
pq.write_table(
table, str(filepath),
compression='snappy',
row_group_size=100000, # 每个行组10万行
use_dictionary=True, # 字典编码(适合 symbol 等低基数列)
write_statistics=True # 写入列统计(min/max/null_count)
)
written_files.append(str(filepath))
return written_files点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:函数 `write_tick_to_parquet`(将 Tick 数据写入分区 Parquet 文件)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。
2.3 ClickHouse 表设计
下面用 Python 把 DDL 管理起来(真实项目里通常也是这么做的:DDL 入版本库,由迁移脚本执行), 顺带打印出这套设计的关键取舍,以及稀疏索引带来的扫描量差异。
此处有展示代码展开 ▼
python
# ClickHouse MergeTree 表定义(适用于 Tick 数据)
DDL_TICK_TRADES = """
CREATE TABLE IF NOT EXISTS tick_trades (
timestamp DateTime64(9, 'UTC'),
symbol LowCardinality(String),
exchange LowCardinality(String),
price Float64,
size UInt64,
trade_id UInt64,
trade_condition LowCardinality(String),
bid_price Float64,
ask_price Float64,
insert_time DateTime DEFAULT now()
)
ENGINE = MergeTree()
PARTITION BY toYYYYMMDD(timestamp) -- 按日分区
ORDER BY (symbol, timestamp) -- 排序键:先 symbol 后时间
TTL timestamp + INTERVAL 30 DAY -- 30天自动过期(按需调整)
SETTINGS index_granularity = 8192; -- 稀疏索引粒度
"""
# 物化视图:预聚合 1 分钟 K 线(写入 tick_trades 时增量维护)
DDL_MINUTE_BARS = """
CREATE MATERIALIZED VIEW tick_minute_bars
ENGINE = SummingMergeTree()
PARTITION BY toYYYYMMDD(timestamp)
ORDER BY (symbol, timestamp)
AS SELECT
symbol,
toStartOfMinute(timestamp) AS timestamp,
argMin(price, timestamp) AS open,
max(price) AS high,
min(price) AS low,
argMax(price, timestamp) AS close,
sum(size) AS volume,
count() AS trade_count
FROM tick_trades
GROUP BY symbol, timestamp;
"""
MIGRATIONS = [('tick_trades', DDL_TICK_TRADES), ('tick_minute_bars', DDL_MINUTE_BARS)]
def apply_migrations(client=None):
"""真实环境传入 clickhouse_driver.Client;这里没有连接就只做 dry-run。"""
for name, ddl in MIGRATIONS:
stmt = ddl.strip()
if client is None:
print(f"[dry-run] would execute DDL for {name} ({len(stmt)} chars, "
f"{len(stmt.splitlines())} lines)")
else:
client.execute(stmt)
print(f"[applied] {name}")
apply_migrations() # 未连接 ClickHouse → dry-run
print("\n=== 设计要点 ===")
DESIGN_NOTES = [
('ENGINE', 'MergeTree', '后台归并有序数据块,写入吞吐高'),
('PARTITION BY', 'toYYYYMMDD(timestamp)', '按日分区 → 时间范围查询直接裁掉无关分区'),
('ORDER BY', '(symbol, timestamp)', '先 symbol 后时间:单票时间序列扫描连续'),
('LowCardinality', 'symbol / exchange', '字典编码,重复字符串只存一次'),
('TTL', 'timestamp + 30 DAY', '热数据留 30 天,过期自动落冷/删除'),
('index_granularity', '8192', '稀疏索引:每 8192 行一个标记'),
]
for k, v, why in DESIGN_NOTES:
print(f" {k:<18}{v:<26}{why}")
# 稀疏索引 + 分区裁剪的效果:查 1 只票 1 天的 tick
TOTAL_ROWS = 5000 * 252 * 20_000 # 全市场一年成交流水行数
ROWS_PER_SYMBOL_DAY = 20_000
GRANULARITY = 8192
partition_rows = 5000 * ROWS_PER_SYMBOL_DAY # 单日分区行数
granules_scanned = max(1, ROWS_PER_SYMBOL_DAY // GRANULARITY + 1)
print("\n=== 查询 “某只票某一天的全部 tick” 的扫描量 ===")
print(f" 全表行数: {TOTAL_ROWS:>15,}")
print(f" 分区裁剪后: {partition_rows:>15,} ({partition_rows / TOTAL_ROWS:.4%})")
print(f" 稀疏索引命中 granule: {granules_scanned:>15,} 个 "
f"→ 实际读取约 {granules_scanned * GRANULARITY:,} 行")
print(f" 相对全表扫描节省: {1 - granules_scanned * GRANULARITY / TOTAL_ROWS:.6%}")
print("\n=== 排序键顺序的影响 ===")
print(" ORDER BY (symbol, timestamp) → 单票回放快;跨票同一时刻切片慢")
print(" ORDER BY (timestamp, symbol) → 全市场快照快;单票回放要跳着读")
print(" 两类查询都重要时,用第二张表或 PROJECTION 各存一份排序")点击展开可浏览运行结果
[dry-run] would execute DDL for tick_trades (686 chars, 17 lines) [dry-run] would execute DDL for tick_minute_bars (477 chars, 15 lines) === 设计要点 === ENGINE MergeTree 后台归并有序数据块,写入吞吐高 PARTITION BY toYYYYMMDD(timestamp) 按日分区 → 时间范围查询直接裁掉无关分区 ORDER BY (symbol, timestamp) 先 symbol 后时间:单票时间序列扫描连续 LowCardinality symbol / exchange 字典编码,重复字符串只存一次 TTL timestamp + 30 DAY 热数据留 30 天,过期自动落冷/删除 index_granularity 8192 稀疏索引:每 8192 行一个标记 === 查询 “某只票某一天的全部 tick” 的扫描量 === 全表行数: 25,200,000,000 分区裁剪后: 100,000,000 (0.3968%) 稀疏索引命中 granule: 3 个 → 实际读取约 24,576 行 相对全表扫描节省: 99.999902% === 排序键顺序的影响 === ORDER BY (symbol, timestamp) → 单票回放快;跨票同一时刻切片慢 ORDER BY (timestamp, symbol) → 全市场快照快;单票回放要跳着读 两类查询都重要时,用第二张表或 PROJECTION 各存一份排序
三、数据压缩与索引策略
3.1 压缩算法选择
此处有展示代码展开 ▼
python
def tick_data_compression_benchmark(sample_data: pd.DataFrame) -> dict:
"""
比较不同压缩算法对 Tick 数据的压缩效果。
Tick 数据的特殊性:
- timestamp 列具有高局部性(delta 编码极有效)
- price 列增量通常很小(适合 delta + 熵编码)
- size 列分布不均
- symbol 列低基数(字典编码极有效)
"""
import time
import io
results = {}
compressions = ['snappy', 'gzip', 'zstd', 'lz4', 'brotli']
for comp in compressions:
buf = io.BytesIO()
start = time.time()
table = pa.Table.from_pandas(sample_data)
pq.write_table(table, buf,
compression=comp,
use_dictionary=True)
compressed_size = buf.tell()
original_size = sample_data.memory_usage(deep=True).sum()
ratio = original_size / compressed_size
elapsed = time.time() - start
results[comp] = {
'compressed_bytes': compressed_size,
'original_bytes': original_size,
'compression_ratio': ratio,
'write_time_seconds': elapsed
}
return results点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:函数 `tick_data_compression_benchmark`(比较不同压缩算法对 Tick 数据的压缩效果)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。
3.2 索引设计
此处有展示代码展开 ▼
python
class TickDataIndex:
"""
Tick 数据的多层索引设计。
索引层级:
1. 分区层:按日期和 symbol 进行文件系统分区(避免扫描无关文件)
2. 文件内元数据:Parquet 的 footer 中存储 min/max 统计
3. 内存缓存:热点数据的内存映射
"""
@staticmethod
def build_symbol_date_catalog(base_path: str) -> pd.DataFrame:
"""
构建 symbol-date 目录索引,加速查询路由。
"""
base = Path(base_path)
entries = []
for symbol_dir in base.iterdir():
if symbol_dir.is_dir():
symbol = symbol_dir.name
for date_dir in symbol_dir.iterdir():
if date_dir.is_dir():
date_str = date_dir.name
parquet_files = list(date_dir.glob('*.parquet'))
total_size = sum(f.stat().st_size
for f in parquet_files)
entries.append({
'symbol': symbol,
'date': date_str,
'file_count': len(parquet_files),
'total_size_bytes': total_size,
'path': str(date_dir)
})
catalog = pd.DataFrame(entries)
catalog['date'] = pd.to_datetime(catalog['date'])
catalog = catalog.sort_values(['symbol', 'date'])
return catalog
@staticmethod
def query_plan(catalog: pd.DataFrame,
symbols: list,
start_date: str,
end_date: str) -> list:
"""
给定查询条件,返回需要读取的文件列表(避免全表扫描)。
"""
mask = (
catalog['symbol'].isin(symbols) &
(catalog['date'] >= start_date) &
(catalog['date'] <= end_date)
)
relevant = catalog[mask]
files_to_read = []
for _, row in relevant.iterrows():
for f in Path(row['path']).glob('*.parquet'):
files_to_read.append(str(f))
return files_to_read点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:类 `TickDataIndex`(Tick 数据的多层索引设计)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。
四、查询优化与分片
4.1 常见查询模式及优化
此处有展示代码展开 ▼
python
def tick_query_optimizer(query_type: str,
params: dict) -> str:
"""
根据查询类型和参数生成优化后的查询策略。
Tick数据的典型查询模式:
1. 时间范围查询(某一天某只股票的所有tick)
2. 快照查询(某一时刻全市场的状态)
3. 聚合查询(计算分钟/小时K线)
4. 模式查询(查找特定模式,如大单、快速价格变动)
"""
strategies = {
'time_range_single_symbol': {
'partition_pruning': True,
'min_max_filter': True,
'parallel_reads': 1, # 单品种线形读取足够
'use_precomputed_bars': False
},
'market_snapshot': {
'partition_pruning': True,
'min_max_filter': True,
'parallel_reads': 8, # 多品种并行读取
'use_precomputed_bars': False
},
'aggregation_all_symbols': {
'partition_pruning': True,
'min_max_filter': False, # 聚合不需要精确过滤
'parallel_reads': 16, # 高度并行
'use_precomputed_bars': True # 用物化视图加速
},
'pattern_detection': {
'partition_pruning': True,
'min_max_filter': True,
'parallel_reads': 4,
'use_precomputed_bars': False,
'push_predicate_to_read': True # 谓词下推
}
}
return strategies.get(query_type, strategies['time_range_single_symbol'])点击展开可浏览运行结果
📘 本段代码定义了 1 个函数/类:函数 `tick_query_optimizer`(根据查询类型和参数生成优化后的查询策略)。该片段为教学展示(未包含独立运行的输入数据),可在实战练习中结合真实数据调用。
4.2 分片策略
对于超大规模 Tick 数据库,合理的分片策略是保证查询性能的关键:
- 水平分片:按 symbol 哈希或按时间范围(年度/季度)分片到不同节点
- 热-温-冷数据分层:
- 热数据(最近7天):SSD + 内存缓存
- 温数据(最近3个月):SSD
- 冷数据(历史):HDD + 高压缩比
好的数据基础设施是量化策略的根基。Tick 数据库设计直接影响研究效率——一个精心设计的列式存储可以在数秒内完成一个月的全市场 Tick 分析,而糟糕的设计可能需要数小时甚至根本不可行。有了可靠的数据存储之后,下一步是构建实时 ETL 管道来持续地将原始数据转化为研究就绪的标准化数据。