저는 지난 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 시스템의 심장입니다.

왜 멀티 거래소 통합 스키마가 필수인가

하지만 "11개 API를 그냥 동시에 호출한다"는 접근은 즉시 3가지 문제를 만듭니다:

  1. 각 거래소가 서로 다른 side 컨벤션(buy/sell vs m=true/false vs direction=long/short)을 씀
  2. symbol이 BTCUSDT, BTC-USDT, XBT/USDT, btcusdt 등 11가지 형태로 분산
  3. 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가지를 절대 원칙으로 삼습니다.

  1. Single source of truth (canonical form): 도메인 로직은 오직 정규화된 스키마만 다룹니다. 거래소 원본은 adapter 출구에서 죽입니다.
  2. 이벤트 소싱(event sourcing): 상태가 아닌 이벤트를 저장. 재현 가능, audit-friendly, 결정론적 backtest.
  3. 3중 timestamp: exchange_ts_ms(거래소 보고 시각), receive_ts_ms(로컬 수신), monotonic_ns(로직 처리 순서). 어느 하나만 쓰면 wall clock drift에 당합니다.
  4. 불변성(immutability): Pydantic frozen=True + Decimal 사용. float는 0.1+0.2 같은 오차로 PnL을 깨뜨립니다.
  5. 결정론적 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.Queuemaxsize를 걸고, 특정 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