암호화폐 시장 분석을 하려면 여러 거래소의 데이터를 한 곳에 모아 정규화해야 합니다. 저는 최근에 Tardis, Binance, OKX의 L2 호가창·체결·펀딩비 데이터를 단일 스키마로 묶어 Parquet 파일로 저장하는 파이프라인을 구축했는데, 이 과정에서 HolySheep AI의 LLM API를 활용해 스키마 매핑 코드를 자동 생성하고 검증했습니다. 이 글에서는 실제 운영 환경에서 검증한 코드, 비용 비교, 자주 마주치는 오류 해결법까지 모두 공유합니다.

왜 거래소 데이터를 통합해야 하는가

저는 2024년부터 멀티 거래소 차익거래 봇을 운영하면서 데이터 불일치 문제를 직접 겪었습니다. 같은 시점의 BTC 가격을 Binance는 USDT로, OKX는 USDT로, Tardis는 USD로 표기하기 때문에 단순 결합이 불가능합니다. 또한 각 거래소의 필드명(price vs px vs last), 타임스탬프 정밀도(밀리초 vs 마이크로초), 호가 깊이(20단계 vs 400단계)도 제각각입니다.

해결책은 다음과 같은 통합 스키마를 정의하고, 모든 원본 데이터를 이 형태로 변환한 뒤 Parquet 컬럼형 포맷으로 저장하는 것입니다. Parquet는 압축률이 높고(원본 대비 약 10배), DuckDB·Polars·Spark에서 모두 네이티브로 읽을 수 있어 후속 분석이 매우 빠릅니다.

통합 스키마 정의 (Unified Schema)

세 거래소의 호가창 스냅샷을 하나로 합치기 위한 표준 스키마는 다음과 같습니다. UTC 기준 마이크로초 정밀도 타임스탬프를 기준으로 통일합니다.

"""
unified_schema.py
다중 거래소 호가창 통합 스키마 정의
"""

import pyarrow as pa

통합 호가창 스키마

UNIFIED_ORDERBOOK_SCHEMA = pa.schema([ pa.field("exchange", pa.string(), nullable=False), # binance, okx, tardis pa.field("symbol", pa.string(), nullable=False), # BTC-USDT, BTC-USD pa.field("ts_event_us", pa.int64(), nullable=False), # 이벤트 시각 (마이크로초) pa.field("ts_recv_us", pa.int64(), nullable=False), # 수신 시각 (마이크로초) pa.field("side", pa.string(), nullable=False), # bid / ask pa.field("level", pa.int32(), nullable=False), # 1~400 pa.field("price", pa.float64(), nullable=False), pa.field("size", pa.float64(), nullable=False), pa.field("quote_currency", pa.string(), nullable=False), # USDT / USD ])

통합 체결 스키마

UNIFIED_TRADE_SCHEMA = pa.schema([ pa.field("exchange", pa.string(), nullable=False), pa.field("symbol", pa.string(), nullable=False), pa.field("ts_event_us", pa.int64(), nullable=False), pa.field("trade_id", pa.string(), nullable=False), pa.field("price", pa.float64(), nullable=False), pa.field("size", pa.float64(), nullable=False), pa.field("side", pa.string(), nullable=False), # buy / sell pa.field("fee", pa.float64(), nullable=True), pa.field("fee_currency", pa.string(), nullable=True), ]) UNIFIED_SCHEMAS = { "orderbook": UNIFIED_ORDERBOOK_SCHEMA, "trade": UNIFIED_TRADE_SCHEMA, }

Binance + OKX 실시간 수집 → 통합 → Parquet 저장

Binance WebSocket과 OKX WebSocket을 동시에 구독해 통합 스키마로 변환한 뒤, PyArrow를 통해 분 단위 Parquet 파일로 저장하는 코드입니다. 저는 이 패턴으로 일 평균 약 4GB의 호가창 데이터를 안정적으로 수집하고 있습니다.

"""
collector_binance_okx.py
Binance / OKX 호가창 + 체결 → 통합 Parquet 저장
"""

import asyncio
import json
import time
from datetime import datetime, timezone

import pyarrow as pa
import pyarrow.parquet as pq
import websockets
from collections import defaultdict

SYMBOL = "BTC-USDT"   # OKX 표기 (Binance는 BTCUSDT로 매핑)
DEPTH = 20
OUT_DIR = "./data/parquet"

호가 누적 버퍼 (exchange, symbol, ts_event_us, side) → list of (level, price, size)

orderbook_buffer = defaultdict(list) trade_buffer = defaultdict(list) def binance_to_unified(symbol_binance: str) -> str: """Binance 표기(BTCUSDT) → 통합 표기(BTC-USDT)""" if symbol_binance.endswith("USDT"): return f"{symbol_binance[:-4]}-USDT" return symbol_binance async def binance_stream(): url = f"wss://stream.binance.com:9443/ws/{SYMBOL.replace('-','').lower()}@depth{DEPTH}@100ms" async with websockets.connect(url, ping_interval=20) as ws: while True: raw = json.loads(await ws.recv()) ts_us = int(raw.get("E", time.time() * 1000) *