私が2025年に参加した東京のあるクオントヘッジファンドでは、BTC/USDTが15分間で−7.2%急落した日に、3つの取引所から合計1,187,432件の強制清算(Liquidation)イベントがバースト的に流入しました。生データのままでは約31%がフィールド欠損や単位不整合で分析不能になり、当時のAI分析チームが「正解ラベルの再構築だけで48時間かかった」と苦労していました。本記事では、私が構築したWebSocket → 正規化 → Parquet保存までのパイプラインを、コード付きで完全公開します。
1. 3取引所の強制清算データフォーマット比較
| 項目 | Binance (USDⓈ-M) | OKX (Swap) | Bybit (USDT Perpetual) |
|---|---|---|---|
| WebSocketエンドポイント | wss://fstream.binance.com/ws/!forceOrder@arr | wss://ws.okx.com:8443/ws/v5/public | wss://stream.bybit.com/v5/public/linear |
| トピック | !forceOrder@arr | liquidation-contracts | allLiquidation |
| 価格フィールド | p (string) | px (string) | price (string) |
| 数量フィールド | q (string) | sz (string) | size (string) |
| サイド | S (SELL=ロング清算) | side (buy/sell) | side (Buy/Sell) |
| タイムスタンプ | T (ms epoch) | ts (ms epoch) | timestamp (ms epoch) |
| シンボル | s (BTCUSDT) | instId (BTC-USDT-SWAP) | symbol (BTCUSDT) |
| 履歴取得 | fapi/v1/forceOrders | /api/v5/public/liquidation-orders | /v5/market/recent-trade (代替) |
| レート制限 | 1200 req/min | 20 req/2s | 600 req/5s |
| 再接続挙動 | 自動再送あり | 切断後手動再購読必須 | 5秒以内自動再接続 |
この表からも分かる通り、同じ「強制清算」でも価格・数量・サイド・シンボルの命名と単位が三者三様です。これを統一スキーマに変換するのが本記事のゴールです。
2. 統一スキーマ定義
from dataclasses import dataclass
from typing import Literal
@dataclass(frozen=True)
class LiquidationEvent:
exchange: Literal["binance", "okx", "bybit"]
symbol: str # 統一表記: "BTCUSDT"
side: Literal["long_liquidation", "short_liquidation"]
price: float # USDT建て
quantity: float # ベース通貨建て (BTC, ETH 等)
notional_usdt: float # price * quantity
timestamp_ms: int # 取引所のイベント発生時刻 (UTC)
received_ms: int # 自前で受信した時刻 (latency 計測用)
raw: dict # 元 JSON を保持 (監査用)
3. 3取引所対応の正規化パイプライン実装
私が本番で使っている実装を簡略化して共有します。再接続・指数バックオフ・PING/PONG も込みで700行ほどに収まる設計です。
import asyncio, json, time, logging, websockets
from collections import defaultdict
ENDPOINTS = {
"binance": ("wss://fstream.binance.com/ws/!forceOrder@arr", None),
"okx": ("wss://ws.okx.com:8443/ws/v5/public",
{"op":"subscribe","args":[{"channel":"liquidation-contracts","instType":"SWAP"}]}),
"bybit": ("wss://stream.bybit.com/v5/public/linear",
{"op":"subscribe","args":["allLiquidation.X"]}),
}
def normalize(ex: str, msg: dict):
now = int(time.time() * 1000)
if ex == "binance":
o = msg["o"]
return LiquidationEvent(
exchange="binance", symbol=o["s"],
side="long_liquidation" if o["S"] == "SELL" else "short_liquidation",
price=float(o["p"]), quantity=float(o["q"]),
notional_usdt=float(o["p"]) * float(o["q"]),
timestamp_ms=int(o["T"]), received_ms=now, raw=msg,
)
if ex == "okx":
d = msg["data"][0]
sym = d["instId"].split("-")[0] + d["instId"].split("-")[1] # BTC-USDT-SWAP -> BTCUSDT
return LiquidationEvent(
exchange="okx", symbol=sym,
side="long_liquidation" if d["side"] == "sell" else "short_liquidation",
price=float(d["px"]), quantity=float(d["sz"]),
notional_usdt=float(d["px"]) * float(d["sz"]),
timestamp_ms=int(d["ts"]), received_ms=now, raw=msg,
)
if ex == "bybit":
d = msg["data"]
return LiquidationEvent(
exchange="bybit", symbol=d["symbol"],
side="long_liquidation" if d["side"] == "Sell" else "short_liquidation",
price=float(d["price"]), quantity=float(d["size"]),
notional_usdt=float(d["price"]) * float(d["size"]),
timestamp_ms=int(d["timestamp"]), received_ms=now, raw=msg,
)
async def stream(ex_name: str, out_q: asyncio.Queue):
url, sub = ENDPOINTS[ex_name]
backoff = 1
while True:
try:
async with websockets.connect(url, ping_interval=20) as ws:
if sub: await ws.send(json.dumps(sub))
backoff = 1
async for raw in ws:
msg = json.loads(raw)
ev = normalize(ex_name, msg)
if ev and ev.price >