私が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@arrwss://ws.okx.com:8443/ws/v5/publicwss://stream.bybit.com/v5/public/linear
トピック!forceOrder@arrliquidation-contractsallLiquidation
価格フィールド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/min20 req/2s600 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 >