ある日、私は凌晨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