基于CCXT库的多交易所流式股票数据高效存储方案咨询
问题分析
你的代码中每次获取订单簿快照就立即写入Parquet文件,频繁的小批量写入会导致Parquet文件的元数据不断膨胀,且单个文件体积过大后,每次追加的IO和元数据更新开销急剧增加,最终引发延迟。此外,当前的DataFrame结构将每个价格档位作为一行,进一步放大了写入行数和元数据开销。
下面是针对高效存储和最小化体积的解决方案:
一、优化现有Parquet存储(低成本改造)
1. 批量写入,减少IO次数
不要每次快照都写入,攒够一定数量的快照(比如100次)或固定时间间隔(比如10秒)再批量写入,大幅降低文件操作次数。
async def watch_book(exchange, ticker): last = None limits = {'Binance': 1000, 'Huobi': 150} columns = ['bids_price', 'bids_value', 'asks_price', 'asks_value', 'time'] batch_data = [] # 批量缓存 batch_size = 100 # 每100条快照写入一次 while True: try: orderbook = await exchange.watch_order_book(ticker) bids = {str(float(b[0])): float(b[1]) for b in orderbook['bids']} asks = {str(float(a[0])): float(a[1]) for a in orderbook['asks']} x = np.full(limits[exchange.name], orderbook['timestamp']) data_table = pd.DataFrame( [bids.keys(), bids.values(), asks.keys(), asks.values(), x], index=columns ).T.dropna(axis=1, how='all').replace({None: np.nan}).fillna(0).astype(float).set_index('time') batch_data.append(data_table) # 达到批量阈值时写入 if len(batch_data) >= batch_size: combined = pd.concat(batch_data) file_path = f'{exchange.name}_snap.parquet.gzip' mode = 'append' if os.path.exists(file_path) else 'w' combined.to_parquet( file_path, engine='fastparquet', mode=mode, compression='lz4', # 替换为lz4,比gzip写入快 write_index=True ) batch_data = [] # 清空缓存 except Exception as e: print(f'{exchange.name} failed {type(e)} {e}')
2. 按时间分区存储
避免单个文件无限膨胀,按小时/天生成独立的Parquet文件,比如Binance_snap_2024052014.parquet.gzip,每个文件只存储对应时间段的数据:
# 在批量写入时,按当前时间生成分区文件名 current_hour = datetime.datetime.now().strftime('%Y%m%d%H') file_path = f'{exchange.name}_snap_{current_hour}.parquet.gzip'
3. 调整压缩算法与DataFrame结构
- 压缩算法:用
lz4或zstd替代gzip,在保证不错压缩率的前提下,写入速度提升3-5倍。 - 优化DataFrame结构:将整个订单簿快照作为一行存储(用数组类型列),减少行数和元数据开销:
# 改为每行存储一个完整快照 def format_snapshot(orderbook): bids = [(float(b[0]), float(b[1])) for b in orderbook['bids']] asks = [(float(a[0]), float(a[1])) for a in orderbook['asks']] return pd.DataFrame({ 'time': [orderbook['timestamp']], 'bids': [bids], 'asks': [asks] }).set_index('time') # 写入时直接存储数组列,fastparquet支持复杂类型 combined.to_parquet(..., engine='fastparquet', compression='zstd')
二、替代存储方案
1. Apache Arrow IPC(Feather v2)
Feather v2基于Arrow格式,支持高效追加写入,写入速度比Parquet快2-3倍,压缩率接近Parquet。适合流式数据的低延迟写入:
# 需要安装pyarrow:pip install pyarrow import pyarrow.feather as feather import pyarrow as pa # 批量写入示例 if len(batch_data) >= batch_size: combined = pd.concat(batch_data) table = pa.Table.from_pandas(combined) file_path = f'{exchange.name}_snap.feather' if os.path.exists(file_path): with feather.FeatherWriter(file_path, mode='append') as writer: writer.write_table(table) else: feather.write_feather(combined, file_path, compression='zstd') batch_data = []
2. RocksDB(极致写入速度)
RocksDB是基于LSM树的键值存储,适合高频写入场景,支持自动压缩,写入延迟极低。可以将时间戳作为key,序列化后的订单簿作为value:
# 需要安装rocksdb:pip install python-rocksdb import rocksdb import msgpack db = rocksdb.DB(f'{exchange.name}_orderbook.db', rocksdb.Options(create_if_missing=True)) # 写入时序列化 snapshot = { 'time': orderbook['timestamp'], 'bids': bids, 'asks': asks } db.put(str(orderbook['timestamp']).encode(), msgpack.packb(snapshot)) # 读取时反序列化 data = db.get(b'1716234000000') snapshot = msgpack.unpackb(data)
3. ClickHouse本地表(分析+存储兼顾)
ClickHouse是列式数据库,专为时序数据优化,写入速度极快,压缩率远超Parquet。本地表适合需要后续分析的场景:
-- 先创建本地表 CREATE TABLE orderbook_snap ( time DateTime64(3), exchange String, bids Array(Tuple(Float64, Float64)), asks Array(Tuple(Float64, Float64)) ) ENGINE = MergeTree ORDER BY time PARTITION BY toDate(time) SETTINGS index_granularity = 8192;
# 用clickhouse-connect写入 from clickhouse_connect import get_client client = get_client(host='localhost', port=8123) # 批量写入数据 rows = [(snapshot['time'], exchange.name, snapshot['bids'], snapshot['asks']) for snapshot in batch_data] client.insert('orderbook_snap', rows, column_names=['time', 'exchange', 'bids', 'asks'])
三、最优方案选择
- 兼顾写入速度与分析需求:优先选择优化后的Parquet(批量+分区+zstd压缩)或Feather v2,生态成熟,后续分析工具支持好。
- 极致低延迟写入:选择RocksDB,适合高频流式数据的实时存储。
- 最小化存储体积:选择ClickHouse本地表或ORC格式(zstd压缩),压缩率可达10:1甚至更高。
内容的提问来源于stack exchange,提问作者Ruslan Kirsanov
相关产品推荐
相关产品推荐

