高频数据分析 — 处理逐笔市场数据
加密货币高频交易以毫秒/微秒级时间尺度运行。每日处理数百万笔逐笔数据、应对突发数据流并提取微观结构信号,是专业交易者的核心竞争力。本指南详解支撑Smart Money API每5分钟更新周期的逐笔数据管理与实时分析技术。
核心洞察: 巨鲸不会瞬间移动市场。其大宗交易会在订单簿失衡、买卖价差和吃单量比率中留下可检测特征。Smart Money API实时捕捉这些特征信号。
逐笔数据采集与存储
数据来源
加密货币逐笔数据主要来自:
- 交易所API: Bybit, Binance, Hyperliquid (WebSocket流)
- 数据聚合商: CoinGecko, Kaiko, Tardis.dev (历史+实时)
- 订单簿快照: 每100ms或每次更新时
使用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 = []
订单簿分析
Level 2订单簿快照
每100-500ms捕获完整订单簿:
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 = 看涨(更多买入压力)
返回 买盘量 / 卖盘量 如果 卖盘量 > 0 否则 1.0
定义 获取价差(自身):
最佳买价 = 最大值(自身.买盘.keys())
最佳卖价 = 最小值(自身.卖盘.keys())
返回 (最佳卖价 - 最佳买价) / 最佳买价 # 百分比价差
微观结构信号
从订单簿结构中提取可操作的信号:
- 买盘卖盘失衡: 前10档的买入与卖出量比率
- 价差压缩: 价差收窄 = 强烈信心
- 冰山订单识别: 检测同一价格的部分成交(隐藏量)
- 闪崩: 微秒内突然下跌5%+(流动性事件)
市场微观结构指标
成交量加权平均价格 (VWAP)
比简单收盘价更好的执行基准:
Python — VWAP 计算
定义 计算_vwap(ticks):
# ticks: (价格, 成交量) 元组列表
分子 = 求和(价格 * 成交量 对于 价格, 成交量 在 ticks)
分母 = 求和(成交量 对于 _, 成交量 在 ticks)
返回 分子 / 分母
主动买入/卖出比率
识别哪一方更为激进:
Python — 主动交易分析
定义 获取_主动交易方向(tick):
# 如果交易价格 = 买价,则卖方激进(供应)
# 如果交易价格 = 卖价,则买方激进(需求)
如果 abs(tick.price - 最佳买价) < tick.price_step:
返回 "卖出"
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)