ある日、私は凌晨3時の job を回していて、突如として次のような例外に遭遇しました。

requests.exceptions.ConnectionError: HTTPSConnectionPool(host='api.binance.com', port=443):
Max retries exceeded with url: /api/v3/trades?symbol=BTCUSDT
(Caused by NewConnectionError('<urllib3.connection.HTTPSConnection object>:
Failed to establish a new connection: Connection timed out after 10 seconds'))

別の日には、API key を間違えて投入してしまい、こんなエラーが出ました。

binance.exceptions.ClientError: 401 Unauthorized
IP=None, total=0, used=0, status=INVALID_API_KEY_OR_SIGNATURE

Binance の後に OKX を、Tardis から過去データを、と試していくうちに、各取引所の「フィールド名・データ型・timestamp 単位・約定の方向表記」がバラバラで、DataFrame にそのまま縦結合しても dtype が崩壊する——これが私が実際に踏んだ最初の地雷でした。本記事では、私が 3 ヶ月の運用でたどり着いた「統一スキーマ → Parquet 書き出し」の設計と、現場で出たエラーの対処法を共有します。

なぜ「統一スキーマ」が必要なのか

Binance の /api/v3/trades{ "price": "43250.10", "qty": "0.012", "time": 1716000000000 } を返します。一方、OKX の /api/v5/market/trades{ "px": "43250.1", "sz": "0.012", "ts": "1716000000123", "side": "buy" } を返します。Tardis は CSV/Parquet で正規化されていますが、フィールド名は price / amount、timestamp はナノ秒精度 1716000000123456789 です。

この差を放置すると、Hive partition で price FLOAT64 のはずが、別パーティションでは STRING になって、Spark が次のエラーを吐きます。

org.apache.spark.sql.analysis.DatasetPath: Parquet type mismatch
Expected: DoubleType, Found: BINARY

だからこそ、データを Parquet に「書く前」に必ず正規化するのが、私の運用ルールです。

私が採用した統一スキーマ(最終形)

from dataclasses import dataclass
from typing import Literal

Side = Literal["buy", "sell"]

@dataclass(frozen=True)
class UnifiedTrade:
    exchange: str          # 'binance' | 'okx' | 'tardis'
    symbol: str            # 正規化済 'BTC-USDT'
    ts_ms: int             # epoch milliseconds (UTC)
    side: Side             # 'buy' | 'ell' (taker side)
    price: float           # quote currency per base
    amount: float          # base currency
    trade_id: str          # 取引所横断で一意
    received_ms: int       # ローカル受信時刻(遅延計測用)

ポイントは ts_ms を必ず 1 カラムに固定することです。Tardis のナノ秒は 1000 で割り、Binance のミリ秒はそのまま、OKX の文字列ミリ秒は int() でキャストしてから入れます。

3 取引所を実際に叩く最小コード

import asyncio
import time
import httpx
import pandas as pd
from typing import AsyncIterator

BINANCE  = "https://api.binance.com"
OKX      = "https://www.okx.com"
TARDIS   = "https://api.tardis.dev/v1"

async def binance_trades(client: httpx.AsyncClient, symbol="BTCUSDT") -> AsyncIterator[dict]:
    r = await client.get(f"{BINANCE}/api/v3/tr