高频数据分析 — 处理逐笔市场数据

加密货币高频交易以毫秒/微秒级时间尺度运行。每日处理数百万笔逐笔数据、应对突发数据流并提取微观结构信号,是专业交易者的核心竞争力。本指南详解支撑Smart Money API每5分钟更新周期的逐笔数据管理与实时分析技术。

核心洞察: 巨鲸不会瞬间移动市场。其大宗交易会在订单簿失衡、买卖价差和吃单量比率中留下可检测特征。Smart Money API实时捕捉这些特征信号。

逐笔数据采集与存储

数据来源

加密货币逐笔数据主要来自:

使用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())
返回 (最佳卖价 - 最佳买价) / 最佳买价 # 百分比价差

微观结构信号

从订单簿结构中提取可操作的信号:

市场微观结构指标

成交量加权平均价格 (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)

为您的分析添加实时信号

Smart Money API聚合了来自3家交易所和250多个鲸鱼钱包的交易数据。使用我们预先计算的微观结构信号来增强您自己的实时分析。

立即免费开始 →