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份快照),后续分析无需重建订单簿。
核心逻辑
- 收到
snapshot后,用字典存储买卖盘(价格为键,数量为值,支持快速更新) - 收到
l2update时,直接更新字典(数量为0则删除对应价格) - 单独线程每秒触发一次快照,用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
相关产品推荐
相关产品推荐

