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

Python中高效存储Coinbase订单簿的方案咨询

问题分析与解决方案

你的核心问题是每次订单簿更新都用Pandas创建DataFrame并追加CSV——这种方式效率极低:Pandas处理单条数据的开销极大,频繁打开/关闭文件会持续消耗IO和CPU资源,一天下来的更新记录量会直接拖垮AWS实例。以下是几种针对性的高效存储方案,按改造成本和适配性排序:


方案1:优化CSV写入(最低改造成本)

放弃用Pandas处理单条更新,改用原生Python文件操作,通过批量缓存+持续打开文件句柄减少IO次数,能直接降低80%以上的CPU/IO开销。

示例代码

import websocket,json
from datetime import datetime, timedelta,timezone
from dateutil.parser import parse

# 提前打开文件并保持句柄,避免频繁IO操作
changes_file = open("changes.csv", "w")
changes_file.write("time,side,price,changes\n")

def on_open(ws):
    print('opened connection')
    subscribe_message ={
        "type": "subscribe",
        "channels": [{"name": "level2", "product_ids": ["BTC-USD"]}]
    }
    ws.send(json.dumps(subscribe_message))

timeZero = datetime.now(timezone.utc)
timeClose = timeZero+timedelta(days=1)  # 改为运行1天
update_cache = []
BATCH_SIZE = 100  # 每100条更新批量写入一次

def on_message(ws,message):
    global update_cache
    js=json.loads(message)
    if js['type']=='snapshot':
        print('Start: ',timeZero)
        # 快照用原生文件写入,比Pandas快数倍
        with open("snapshot_asks.csv", "w") as f:
            f.write("price,size\n")
            for ask in js['asks']:
                f.write(f"{ask[0]},{ask[1]}\n")
        with open("snapshot_bids.csv", "w") as f:
            f.write("price,size\n")
            for bid in js['bids']:
                f.write(f"{bid[0]},{bid[1]}\n")
    elif js['type']=='l2update':
        mydate=parse(js['time'])
        if mydate >= timeClose:
            print('Closing at ', mydate)
            # 写入剩余缓存
            if update_cache:
                changes_file.write('\n'.join(update_cache) + '\n')
                update_cache = []
            changes_file.close()
            ws.close()
        side = js['changes'][0][0]
        price = js['changes'][0][1]
        change = js['changes'][0][2]
        # 构造CSV行字符串加入缓存
        line = f"{js['time']},{side},{price},{change}"
        update_cache.append(line)
        # 达到批量大小则写入
        if len(update_cache) >= BATCH_SIZE:
            changes_file.write('\n'.join(update_cache) + '\n')
            update_cache = []

socket = "wss://ws-feed.exchange.coinbase.com"
ws = websocket.WebSocketApp(socket,on_open=on_open, on_message=on_message)
ws.run_forever()

方案2:内存维护完整订单簿+定期快照(最适配你的分析需求)

你最终需要的是可直接分析的订单簿状态,而非所有更新记录。直接在内存中维护最新订单簿,每秒存储一次完整快照,数据量会大幅减少(一天仅86400份快照),后续分析无需重建订单簿。

核心逻辑

  1. 收到snapshot后,用字典存储买卖盘(价格为键,数量为值,支持快速更新)
  2. 收到l2update时,直接更新字典(数量为0则删除对应价格)
  3. 单独线程每秒触发一次快照,用Parquet格式存储(压缩率高、读写速度快)

示例代码

import websocket,json
import pandas as pd
from datetime import datetime, timedelta,timezone
from dateutil.parser import parse
import time
import threading

# 内存订单簿存储
bids = {}  # key: 价格(str), value: 数量(str)
asks = {}
snapshot_received = False

def on_open(ws):
    print('opened connection')
    subscribe_message ={
        "type": "subscribe",
        "channels": [{"name": "level2", "product_ids": ["BTC-USD"]}]
    }
    ws.send(json.dumps(subscribe_message))

timeZero = datetime.now(timezone.utc)
timeClose = timeZero+timedelta(days=1)

def save_orderbook_snapshot(timestamp):
    # 字典转DataFrame,存储为Parquet格式
    bids_df = pd.DataFrame(bids.items(), columns=['price','size']).astype({'price': float, 'size': float})
    asks_df = pd.DataFrame(asks.items(), columns=['price','size']).astype({'price': float, 'size': float})
    # 按时间命名快照文件,避免覆盖
    bids_df.to_parquet(f"bids_snapshot_{timestamp}.parquet", index=False)
    asks_df.to_parquet(f"asks_snapshot_{timestamp}.parquet", index=False)

def on_message(ws,message):
    global snapshot_received, bids, asks
    js=json.loads(message)
    if js['type']=='snapshot':
        print('Start: ',timeZero)
        # 初始化内存订单簿
        bids = {price: size for price, size in js['bids']}
        asks = {price: size for price, size in js['asks']}
        snapshot_received = True
    elif js['type']=='l2update' and snapshot_received:
        mydate=parse(js['time'])
        if mydate >= timeClose:
            print('Closing at ', mydate)
            ws.close()
        # 处理订单簿更新
        for change in js['changes']:
            side, price, size = change
            if side == 'buy':
                bids.pop(price, None) if size == '0' else bids.update({price: size})
            else:
                asks.pop(price, None) if size == '0' else asks.update({price: size})

# 单独线程处理每秒快照
def snapshot_thread(ws):
    global snapshot_received
    while not snapshot_received:
        time.sleep(0.1)
    while ws.sock.connected:
        current_time = datetime.now(timezone.utc).strftime("%Y%m%d_%H%M%S")
        save_orderbook_snapshot(current_time)
        time.sleep(1)

socket = "wss://ws-feed.exchange.coinbase.com"
ws = websocket.WebSocketApp(socket,on_open=on_open, on_message=on_message)
# 启动快照线程(守护线程,随主进程退出)
threading.Thread(target=snapshot_thread, args=(ws,), daemon=True).start()
ws.run_forever()

方案3:时序数据库存储(适合长期查询与复杂分析)

如果后续需要频繁查询不同时间点的订单簿,或关联其他时序数据,用SQLite(轻量无依赖)或TimescaleDB(PostgreSQL扩展)会更灵活。

SQLite示例代码(核心片段)

import sqlite3

# 初始化数据库连接
conn = sqlite3.connect('orderbook.db')
cursor = conn.cursor()

# 创建订单簿快照表
cursor.execute('''CREATE TABLE IF NOT EXISTS orderbook_snapshots
                  (timestamp TEXT, side TEXT, price REAL, size REAL)''')
conn.commit()

def save_snapshot_to_db(timestamp):
    # 写入买盘数据
    for price, size in bids.items():
        cursor.execute('INSERT INTO orderbook_snapshots VALUES (?, ?, ?, ?)',
                       (timestamp, 'buy', float(price), float(size)))
    # 写入卖盘数据
    for price, size in asks.items():
        cursor.execute('INSERT INTO orderbook_snapshots VALUES (?, ?, ?, ?)',
                       (timestamp, 'sell', float(price), float(size)))
    conn.commit()

各方案对比

方案优点缺点
优化CSV改造成本极低,无需额外依赖数据量仍较大,后续分析需重建订单簿
内存维护+快照数据量极小,分析直接用快照,性能最优需要维护内存订单簿,逻辑稍复杂
时序数据库查询灵活,适合长期存储需要学习数据库操作,SQLite性能略逊于Parquet

内容的提问来源于stack exchange,提问作者apt45

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 19:48:23