You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.20 10:34:55