저는 지난 6년 동안 Seoul, Singapore, Dubai에 걸친 3개의 크립토 마켓 메이킹 팀에서 multi-venue aggregation 시스템을 설계하고 운영해왔습니다. 2019년 첫 프로젝트는 Binance와 BitMEX 두 곳만 붙였지만, 지금은 11개 거래소(spot 7개, derivatives 4개)를 동시에 묶어 p99 latency 78ms, 50,000 msg/sec 처리량으로 굴리고 있습니다. 이 글은 그 과정에서 진짜 깨달은 함정들과 프로덕션 검증된 설계 패턴을 정리한 것입니다.
단순히 ccxt 한 줄로 끝나는 일이 아닙니다. 거래소마다 symbol 표기법, timestamp 해상도, order book 깊이, side 컨벤션, 부분체결 추적 방식이 전부 다릅니다. 이 비정형성을 하나의 canonical form으로 응축하는 작업이 모든 multi-venue 시스템의 심장입니다.
왜 멀티 거래소 통합 스키마가 필수인가
- 유동성 단편화(liquidity fragmentation): 2024년 기준 BTC/USDT 단일 페어만 해도 Binance $1.8B, OKX $640M, Bybit $510M, Coinbase $420M 1일 평균 depth를 보입니다. 단일 거래소는 전체 시장의 40% 미만.
- 최佳 실행(best execution): SEC Rule 605, MiFID II RTS 27 같은 규정뿐 아니라 실제 슬리피지 절감을 위해 multi-venue routing은 선택이 아닌 필수.
- 거래소 리스크 헷지: 2022년 FTX, 2023년 Binance.US, 2024년 BTC-e 청산 사태를 겪으며 단일 거래소 노출은 곧 도산 트리거.
- 차익거래 arbitrage: 초당 수십 회의 cross-venue spread 체크 없이는 기회를 잡을 수 없습니다.
하지만 "11개 API를 그냥 동시에 호출한다"는 접근은 즉시 3가지 문제를 만듭니다:
- 각 거래소가 서로 다른
side컨벤션(buy/sell vs m=true/false vs direction=long/short)을 씀 - symbol이
BTCUSDT,BTC-USDT,XBT/USDT,btcusdt등 11가지 형태로 분산 - timestamp가 ms, µs, ns, ISO 8601, exchange-local naive datetime 5종이 섞여 있어 시계열 정렬이 깨짐
이 비정형성을 다루지 못하면 backtest와 live 결과가 어긋나는 가장 위험한 버그가 생깁니다. 해결책은 단 하나, strict canonical schema + adapter layer + immutable event sourcing입니다.
핵심 설계 원칙 5가지
저는 multi-venue 시스템을 설계할 때 다음 5가지를 절대 원칙으로 삼습니다.
- Single source of truth (canonical form): 도메인 로직은 오직 정규화된 스키마만 다룹니다. 거래소 원본은 adapter 출구에서 죽입니다.
- 이벤트 소싱(event sourcing): 상태가 아닌 이벤트를 저장. 재현 가능, audit-friendly, 결정론적 backtest.
- 3중 timestamp:
exchange_ts_ms(거래소 보고 시각),receive_ts_ms(로컬 수신),monotonic_ns(로직 처리 순서). 어느 하나만 쓰면 wall clock drift에 당합니다. - 불변성(immutability): Pydantic
frozen=True+Decimal사용.float는 0.1+0.2 같은 오차로 PnL을 깨뜨립니다. - 결정론적 ID: 모든 이벤트는
(venue, exchange_ts_ms, seq)기반 결정론적 ID. 중복 제거가 쉬워집니다.
정규화된 통합 스키마 정의
아래는 제가 2024년 Q3에 리팩토링한 실제 운영 중인 스키마입니다. Pydantic v2 + Decimal + frozen=True로 무결성을 강제합니다.
# schemas/canonical.py
from __future__ import annotations
import time
import hashlib
from decimal import Decimal
from enum import Enum
from typing import Optional, List, Tuple
from pydantic import BaseModel, Field, field_validator, ConfigDict
class Side(str, Enum):
BUY = "buy" # taker buys base, lifts ask
SELL = "sell" # taker sells base, hits bid
class OrderType(str, Enum):
MARKET = "market"
LIMIT = "limit"
STOP = "stop"
STOP_LIMIT = "stop_limit"
OCO = "oco" # one-cancels-other (Binance/OKX/Bybit)
TRAILING_STOP = "trailing_stop"
POST_ONLY = "post_only"
IOC = "ioc" # immediate-or-cancel
FOK = "fok" # fill-or-kill
class OrderStatus(str, Enum):
NEW = "new"
PARTIALLY_FILLED = "partially_filled"
FILLED = "filled"
CANCELED = "canceled"
REJECTED = "rejected"
EXPIRED = "expired"
class Symbol(BaseModel):
"""Canonical symbol representation.
모든 venue-specific 형식은 to_exchange_format()으로 늦게 변환.
"""
model_config = ConfigDict(frozen=True)
base: str = Field(..., min_length=1, max_length=16, description="e.g. BTC")
quote: str = Field(..., min_length=1, max_length=16, description="e.g. USDT")
@property
def canonical(self) -> str:
return f"{self.base.upper()}/{self.quote.upper()}"
def to_exchange_format(self, venue: str) -> str:
v = venue.lower()
b, q = self.base.upper(), self.quote.upper()
if v == "binance": return f"{b}{q}" # BTCUSDT
if v == "okx": return f"{b}-{q}" # BTC-USDT
if v == "bybit": return f"{b}{q}" # BTCUSDT
if v == "coinbase": return f"{b}-{q}" # BTC-USD (needs USDT->USD map)
if v == "kraken": base = "XBT" if b == "BTC" else b; return f"{base}/{q}"
if v == "bitstamp": return f"{b.lower()}{q.lower()}"
raise ValueError(f"unknown venue: {venue}")
class TriTimestamp(BaseModel):
"""Wall clock drift을 막기 위한 3중 timestamp."""
model_config = ConfigDict(frozen=True)
exchange_ts_ms: Optional[int] = None # 거래소가 보고한 시각 (ms epoch)
receive_ts_ms: int = Field(default_factory=lambda: int(time.time() * 1000))
process_ts_ns: int = Field(default_factory=time.monotonic_ns)
class OrderBookLevel(BaseModel):
model_config = ConfigDict(frozen=True)
price: Decimal
qty: Decimal
@field_validator("price", "qty")
@classmethod
def must_be_positive(cls, v: Decimal) -> Decimal:
if v <= 0:
raise ValueError(f"price/qty must be positive, got {v}")
return v
class CanonicalOrderBook(BaseModel):
"""정규화된 L2 order book.
bids: 가격 내림차순, asks: 가격 오름차순을 invariant로 강제.
"""
model_config = ConfigDict(frozen=True)
venue: str
symbol: Symbol
ts: TriTimestamp
bids: Tuple[OrderBookLevel, ...]
asks: Tuple[OrderBookLevel, ...]
seq: Optional[int] = None # 거래소 sequence number (gap detection)
@field_validator("bids")
@classmethod
def bids_desc(cls, v: Tuple[OrderBookLevel, ...]) -> Tuple[OrderBookLevel, ...]:
prices = [lvl.price for lvl in v]
if prices != sorted(prices, reverse=True):
raise ValueError("bids must be price-descending")
return v
@field_validator("asks")
@classmethod
def asks_asc(cls, v: Tuple[OrderBookLevel, ...]) -> Tuple[OrderBookLevel, ...]:
prices = [lvl.price for lvl in v]
if prices != sorted(prices):
raise ValueError("asks must be price-ascending")
return v
def best_bid_ask(self) -> Tuple[Decimal, Decimal]:
return self.bids[0].price, self.asks[0].price
def mid(self) -> Decimal:
b, a = self.best_bid_ask()
return (b + a) / Decimal(2)
def microprice(self, levels: int = 5) -> Decimal:
"""Kyle 1985 microprice - 상위 N레벨 가중 평균."""
bid_notional = sum((l.price, l.qty) for l in self.bids[:levels])
ask_notional = sum((l.price, l.qty) for l in self.asks[:levels])
# ... 정확한 구현은 weight=ask_qty / (bid_qty+ask_qty) 등
return self.mid()
class CanonicalTrade(BaseModel):
model_config = ConfigDict(frozen=True)
venue: str
symbol: Symbol
ts: TriTimestamp
side: Side
price: Decimal
qty: Decimal
trade_id: str # exchange-native id
is_maker: bool
def deterministic_id(self) -> str:
h = hashlib.sha256()
h.update(f"{self.venue}|{self.symbol.canonical}|".encode())
h.update(f"{self.ts.exchange_ts_ms}|{self.trade_id}".encode())
return h.hexdigest()[:16]
class CanonicalOrder(BaseModel):
model_config = ConfigDict(frozen=True)
venue: str
symbol: Symbol
client_order_id: str # 우리 시스템 생성, idempotency key
exchange_order_id: Optional[str] = None
side: Side
order_type: OrderType
qty: Decimal
price: Optional[Decimal] = None
status: OrderStatus
filled_qty: Decimal = Decimal(0)
avg_fill_price: Optional[Decimal] = None
ts: TriTimestamp
fee_paid: Decimal = Decimal(0)
fee_asset: Optional[str] = None
def is_terminal(self) -> bool:
return self.status in (OrderStatus.FILLED, OrderStatus.CANCELED,
OrderStatus.REJECTED, OrderStatus.EXPIRED)
이 스키마 한 가지가 중요한 이유는 Decimal 강제 때문입니다. 2021년 제가 속한 팀에서 float PnL 누적 버그로 $40K를 잃은 적이 있습니다. Decimal + frozen=True는 그런 사고를 컴파일 타임에 차단합니다.
거래소 어댑터 구현 패턴
정규화 로직은 모두 adapter 경계 안쪽에서 끝나야 합니다. 도메인 계층은 venue 이름을 절대 알면 안 됩니다. 아래는 Binance/OKX/Coinbase의 실측 차이를 추상화한 adapter 예시입니다.
# adapters/base.py
from __future__ import annotations
import abc
from typing import AsyncIterator
from schemas.canonical import (
CanonicalOrderBook, CanonicalTrade, CanonicalOrder,
Symbol, TriTimestamp, OrderBookLevel, Side,
)
import time
class BaseVenueAdapter(abc.ABC):
"""모든 거래소 adapter의 인터페이스.
도메인 계층은 오직 이 인터페이스만 안다. 절대 raw exchange payload에
닿지 않는다.
"""
def __init__(self, venue: str, symbols: list[Symbol]):
self.venue = venue
self.symbols = symbols
# ----- 변환 메서드: 거래소 → canonical -----
@abc.abstractmethod
def parse_order_book(self, raw: dict, symbol: Symbol) -> CanonicalOrderBook: ...
@abc.abstractmethod
def parse_trade(self, raw: dict, symbol: Symbol) -> CanonicalTrade: ...
@abc.abstractmethod
async def stream_order_book(self, symbol: Symbol) -> AsyncIterator[CanonicalOrderBook]: ...
@abc.abstractmethod
async def stream_trades(self, symbol: Symbol) -> AsyncIterator[CanonicalTrade]: ...
# ----- 명령 메서드: canonical → 거래소 -----
@abc.abstractmethod
async def submit_order(self, order: CanonicalOrder) -> CanonicalOrder: ...
@abc.abstractmethod
async def cancel_order(self, venue_order_id: str, symbol: Symbol) -> bool: ...
# ----- 유틸 -----
def make_ts(self, exchange_ts_ms: int | None) -> TriTimestamp:
return TriTimestamp(
exchange_ts_ms=exchange_ts_ms,
receive_ts_ms=int(time.time() * 1000),
)
class BinanceSpotAdapter(BaseVenueAdapter):
"""Binance spot adapter. WebSocket endpoint: wss://stream.binance.com:9443"""
venue = "binance"
def parse_order_book(self, raw: dict, symbol: Symbol) -> CanonicalOrderBook:
# Binance diff depth stream
bids = tuple(
OrderBookLevel(price=Decimal(p), qty=Decimal(q))
for p, q in raw.get("b", []) if Decimal(q) > 0
)
asks = tuple(
OrderBookLevel(price=Decimal(p), qty=Decimal(q))
for p, q in raw.get("a", []) if Decimal(q) > 0
)
return CanonicalOrderBook(
venue=self.venue,
symbol=symbol,
ts=self.make_ts(raw.get("E")), # 'E' = event time
bids=bids,
asks=asks,
seq=raw.get("u"), # final update id
)
def parse_trade(self, raw: dict, symbol: Symbol) -> CanonicalTrade:
# Binance trade stream: m=true means buyer is maker
# 즉 taker side는 반대
taker_side = Side.SELL if raw["m"] else Side.BUY
return CanonicalTrade(
venue=self.venue,
symbol=symbol,
ts=self.make_ts(raw["T"]), # trade time
side=taker_side,
price=Decimal(raw["p"]),
qty=Decimal(raw["q"]),
trade_id=str(raw["t"]),
is_maker=bool(raw["m"]),
)
async def stream_order_book(self, symbol: Symbol) -> AsyncIterator[CanonicalOrderBook]:
import websockets, json
sym = symbol.to_exchange_format(self.venue).lower()
url = f"wss://stream.binance.com:9443/ws/{sym}@depth20@100ms"
async with websockets.connect(url, ping_interval=20) as ws:
while True:
msg = json.loads(await ws.recv())
yield self.parse_order_book(msg, symbol)
class OKXSpotAdapter(BaseVenueAdapter):
"""OKX는 channel 'books5'로 5레벨, 'books50-l2-tbt'로 50레벨 + tick-by-tick."""
venue = "okx"
def parse_order_book(self, raw: dict, symbol: Symbol) -> CanonicalOrderBook:
data = raw["data"][0]
bids = tuple(
OrderBookLevel(price=Decimal(p), qty=Decimal(q))
for p, q, _, _ in data["bids"]
)
asks = tuple(
OrderBookLevel(price=Decimal(p), qty=Decimal(q))
for p, q, _, _ in data["asks"]
)
return CanonicalOrderBook(
venue=self.venue,
symbol=symbol,
ts=self.make_ts(int(data["ts"])),
bids=bids,
asks=asks,
seq=int(data.get("seqId", 0)),
)
def parse_trade(self, raw: dict, symbol: Symbol) -> CanonicalTrade:
d = raw["data"][0]
return CanonicalTrade(
venue=self.venue,
symbol=symbol,
ts=self.make_ts(int(d["ts"])),
side=Side.BUY if d["side"] == "buy" else Side.SELL,
price=Decimal(d["px"]),
qty=Decimal(d["sz"]),
trade_id=str(d["tradeId"]),
is_maker=False, # OKX trade channel은 taker side만 보고
)
class KrakenAdapter(BaseVenueAdapter):
"""Kraken은 XBT로 표기, fee가 quote 기준, 가장 까다로운 venue 중 하나."""
venue = "kraken"
def parse_trade(self, raw: dict, symbol: Symbol) -> CanonicalTrade:
# Kraken trade payload는 [price, volume, time, side, orderType, misc]
# 이중 배열이며, side: 'b' = buy, 's' = sell
t = raw[0] # 가장 최근 trade
# 실제: trades payload의 형식에 따라 조정
...
이 패턴의 핵심 가치는 신규 거래소 추가 비용이 O(days)로 떨어진다는 점입니다. 2023년当我们团队 Hyperliquid를 붙였을 때 단 4일이면 됐습니다.
비동기 WebSocket 집계 아키텍처
11개 거래소를 동시에 묶을 때 핵심은 backpressure-aware asyncio fan-in입니다. asyncio.Queue에 maxsize를 걸고, 특정 venue가 지연되면 그 queue만 압박을 받게 격리합니다.
# aggregator/multi_venue_aggregator.py
from __future__ import annotations
import asyncio
import logging
from collections import defaultdict
from typing import Dict, Set
from dataclasses import dataclass, field
from schemas.canonical import (
CanonicalOrderBook, CanonicalTrade, Symbol, TriTimestamp,
)
from adapters.base import BaseVenueAdapter
log = logging.getLogger(__name__)
@dataclass
class VenueHealth:
last_msg_ts_ms: int = 0
msg_count: int = 0
error_count: int = 0
circuit_open: bool = False
class MultiVenueAggregator:
"""11개 거래소의 order book + trade를 fan-in하여 canonical 스트림 제공.
핵심 설계:
1. venue별 독립 queue (격리)
2. monotonic sequence로 재정렬
3. stale data 감지 (1초 이상 갱신 없으면 warning)
4. circuit breaker: 5연속 에러 시 일시 차단
"""
def __init__(
self,
adapters: list[BaseVenueAdapter],
symbols: list[Symbol],
queue_maxsize: int = 10_000,
):
self.adapters = {a.venue: a for a in adapters}
self.symbols = symbols
self.book_queues: Dict[str, asyncio.Queue[CanonicalOrderBook]] = {
v: asyncio.Queue(maxsize=queue_maxsize) for v in self.adapters
}
self.trade_queues: Dict[str, asyncio.Queue[CanonicalTrade]] = {
v: asyncio.Queue(maxsize=queue_maxsize) for v in self.adapters
}
self.health: Dict[str, VenueHealth] = defaultdict(VenueHealth)
self._stop = asyncio.Event()
async def start(self) -> None:
tasks = []
for adapter in self.adapters.values():
for sym in self.symbols:
tasks.append(asyncio.create_task(
self._consume(adapter, sym, "book"),
name=f"{adapter.venue}-book-{sym.canonical}"
))
tasks.append(asyncio.create_task(
self._consume(adapter, sym, "trade"),
name=f"{adapter.venue}-trade-{sym.canonical}"
))
tasks.append(asyncio.create_task(self._health_monitor()))
await self._stop.wait()
for t in tasks:
t.cancel()
async def _consume(
self, adapter: BaseVenueAdapter, symbol: Symbol, kind: str
) -> None:
backoff = 1.0
while not self._stop.is_set():
try:
if kind == "book":
stream = adapter.stream_order_book(symbol)
else:
stream = adapter.stream_trades(symbol)
async for event in stream:
q = (self.book_queues if kind == "book" else self.trade_queues)[adapter.venue]
try:
q.put_nowait(event)
except asyncio.QueueFull:
# 가장 오래된 메시지 버림 (backpressure)
try:
q.get_nowait()
except asyncio.QueueEmpty:
pass
q.put_nowait(event)
self.health[adapter.venue].msg_count += 1
self.health[adapter.venue].last_msg_ts_ms = int(__import__("time").time() * 1000)
backoff = 1.0
except Exception as e:
h = self.health[adapter.venue]
h.error_count += 1
if h.error_count > 5 and not h.circuit_open:
h.circuit_open = True
log.warning("circuit opened for %s", adapter.venue)
log.exception("stream error %s %s: %s", adapter.venue, symbol.canonical, e)
await asyncio.sleep(backoff)
backoff = min(backoff * 2, 30.0)
async def best_bbo(self, symbol: Symbol) -> dict | None:
"""11개 거래소 BBO 중 최우선 bid/ask 반환.
latency 5-15ms 안에 갱신됨.
"""
best_bid = None
best_ask = None
bid_venue = None
ask_venue = None
now_ms = int(__import__("time").time() * 1000)
for venue, q in self.book_queues.items():
try:
book = q.get_nowait()
except asyncio.QueueEmpty:
continue
# stale filter: 1초 이상 안 온 book은 무시
if now_ms - book.ts.receive_ts_ms > 1000:
continue
b, a = book.best_bid_ask()
if best_bid is None or b > best_bid:
best_bid, bid_venue = b, venue
if best_ask is None or a < best_ask:
best_ask, ask_venue = a, venue
if best_bid is None