高頻數據分析 — 處理逐筆市場數據
加密貨幣高頻交易以毫秒/微秒級時序運作。每日處理數百萬筆逐筆數據、應對突發數據流並提取微結構信號,正是專業交易者的核心能力。本指南將詳解支撐Smart Money API實現5分鐘刷新週期的逐筆數據管理與即時分析技術。
關鍵洞察: 鯨魚不會瞬間移動市場。其大額交易會在訂單簿失衡、買賣價差及吃單量比率中留下可檢測特徵。Smart Money API即時捕捉這些特徵。
逐筆數據採集與存儲
數據來源
加密貨幣逐筆數據來自:
- 交易所API: Bybit、Binance、Hyperliquid(WebSocket流)
- 聚合平台: CoinGecko、Kaiko、Tardis.dev(歷史+實時)
- 訂單簿快照: 每100毫秒或每次更新時採集
使用Parquet高效存儲逐筆數據
存儲壓縮列式數據以實現快速查詢:
Python — 逐筆數據處理管道
import pandas as pd
import pyarrow.parquet as pq
# 即時逐筆緩衝區(累積1小時後保存)
tick_buffer = []
def on_tick(exchange, symbol, price, size, side, timestamp):
tick_buffer.append({
"timestamp": timestamp,
"exchange": exchange,
"symbol": symbol,
"price": price,
"size": size,
"side": side # "buy" 或 "sell"
})
# 每小時壓縮保存一次
if len(tick_buffer) > 1_000_000:
df = pd.DataFrame(tick_buffer)
pq.write_table(
pa.Table.from_pandas(df),
f"ticks/{symbol}_{timestamp:%Y%m%d_%H}.parquet",
compression="snappy"
)
tick_buffer = []
訂單簿分析
二級訂單簿快照
每100-500毫秒採集完整訂單簿:
Python — 訂單簿處理
class OrderBook:
def __init__(self):
self.bids = {} # 價格 -> 數量
self.asks = {} # 價格 -> 數量
def update(self, side, price, size):
if side == "bid":
if size == 0: del self.bids[price]
else: self.bids[price] = size
else:
if size == 0: del self.asks[price]
else: self.asks[price] = size
def get_imbalance(self, depth=10):
# 獲取最佳10檔買賣單
top_bids = sorted(self.bids.items(), reverse=True)[:depth]
top_asks = sorted(self.asks.items())[:depth]
bid_volume = sum(size for _, size in top_bids)
ask_volume = sum(size for _, size in top_asks)
# 不平衡:>1 = 看漲(更多買壓)
return bid_volume / ask_volume if ask_volume > 0 else 1.0
def get_spread(self):
best_bid = max(self.bids.keys())
best_ask = min(self.asks.keys())
return (best_ask - best_bid) / best_bid # 百分比價差
微結構信號
從訂單簿結構中提取可操作的信號:
- 買賣不平衡: 前10層級的買賣量比率
- 價差壓縮: 價差縮小 = 強烈信心
- 冰山識別: 檢測相同價格的部分成交(隱藏量)
- 閃崩: 微秒內突然下跌5%+(流動性事件)
市場微結構指標
成交量加權平均價格 (VWAP)
比簡單收盤價更好的執行基準:
Python — VWAP 計算
def calculate_vwap(ticks):
# ticks: (價格, 成交量) 元組列表
numerator = sum(price * volume for price, volume in ticks)
denominator = sum(volume for _, volume in ticks)
return numerator / denominator
主動買/賣比率
識別哪一方更積極:
Python — 主動方分析
def get_taker_direction(tick):
# 如果交易價格 = 買價,則賣方更積極(供應)
# 如果交易價格 = 賣價,則買方更積極(需求)
if abs(tick.price - best_bid) < tick.price_step:
return "sell"
elif abs(tick.price - best_ask) < tick.price_step:
return "buy"
else:
return "mid" # 價差內,可能是暗池
buy_volume = sum(t.size for t in ticks if get_taker_direction(t) == "buy")
sell_volume = sum(t.size for t in ticks if get_taker_direction(t) == "sell")
return buy_volume / (buy_volume + sell_volume) # % 買入
實時管道架構
WebSocket 流處理
連接 Bybit/Binance/Hyperliquid WebSocket 獲取實時數據:
Python — 實時流處理器
import asyncio
import websockets
async def connect_bybit_ticks(symbol):
url = f"wss://stream.bybit.com/v5/public/spot"
async with websockets.connect(url) as ws:
# 訂閱交易流
await ws.send(json.dumps({
"op": "subscribe",
"args": [f"publicTrade.{symbol}"]
}))
async for message in ws:
data = json.loads(message)
for trade 於 數據["data"]:
on_tick(
exchange="bybit",
symbol=symbol,
price=float(trade["price"]),
size=float(trade["size"]),
side=trade["side"],
timestamp=int(trade["time"])
)
基於Redis串流的分散式處理
透過消息隊列處理每日數百萬筆報價:
Python — Redis串流處理
import redis
r = redis.Redis(host='localhost', port=6379)
# 生產者:推送報價至串流
def publish_tick(exchange, symbol, tick):
stream_key = f"ticks:{exchange}:{symbol}"
r.xadd(stream_key, {
"price": tick.price,
"size": tick.size,
"side": tick.side,
"ts": tick.timestamp
})
# 消費者群組讀取並追蹤延遲
r.xgroup_create(stream_key, "analytics", id="$", mkstream=True)