HardПрактика13 min

Биржа: matching engine

Проектирование биржевого matching engine: order book, price-time priority, low latency, FIX-протокол, журналирование и репликация

Matching engine -- сердце любой биржи. Именно он сопоставляет заявки покупателей и продавцов, формирует сделки и поддерживает стакан (order book). От его корректности и предсказуемой latency зависит справедливость торгов и доверие участников рынка.

Функциональные требования

  1. Приём заявок (orders) разных типов: limit, market, stop-limit, IOC (immediate-or-cancel), FOK (fill-or-kill), post-only.
  2. Отмена и модификация заявок.
  3. Сопоставление заявок по принципу price-time priority (сначала лучшая цена, при равных ценах -- кто раньше поставил).
  4. Публикация market data: топ-уровня (L1), полного стакана (L2), trade tape (исполненные сделки).
  5. Журналирование всех событий (orders, trades, cancels) в append-only лог -- для аудита и recovery.
  6. Поддержка нескольких инструментов (symbols): BTC/USDT, ETH/USDT, AAPL, и т.д.

Нефункциональные требования

  1. Latency: p99 end-to-end tick-to-trade < 1 ms, медиана < 100 us внутри matching engine.
  2. Throughput: 100K заявок/сек на один инструмент, 1M+ на кластер.
  3. Fairness: детерминированный порядок обработки, никакой параллельности внутри символа.
  4. Durability: zero data loss для принятых заявок (synchronous replication journal).
  5. Availability: 99.995% в торговые часы, failover < 1s.
  6. Auditability: полный replay истории за любой день.

Оценки нагрузки

  • Пусть 1000 активных инструментов, средняя нагрузка 10K orders/sec на инструмент, пик -- 100K.
  • Peak cluster load: 100K * 1000 = 100M orders/sec (теоретически, практически шардируем).
  • Размер события в journal: ~128 байт (order ID, symbol, price, qty, side, timestamp, client ID).
  • Журнал: 100K * 128 = 12.5 MB/s на инструмент, 300 GB/сутки на горячий символ.
  • Stakan depth: 100K-1M уровней цен, 10M открытых заявок на инструмент в пиках.
  • Network: market data multicast -- сотни Mbps на подписчика, aggregated -- десятки Gbps.

High-Level Architecture

                                  ┌─────────────────────┐
  ┌───────────┐  FIX/WebSocket     │  Gateway Cluster    │
  │  Client   │ ─────────────────> │  (auth, risk check, │
  │  (trader) │ <─ acks/marketdata │   rate limiting)    │
  └───────────┘                    └──────────┬──────────┘
                                              │ sequenced
                                              ▼
                                   ┌─────────────────────┐
                                   │  Sequencer          │
                                   │  (assign global seq │
                                   │   per symbol)       │
                                   └──────────┬──────────┘
                                              │ shard by symbol
                    ┌─────────────────────────┼─────────────────────────┐
                    ▼                         ▼                         ▼
           ┌──────────────────┐      ┌──────────────────┐     ┌──────────────────┐
           │  Matching Engine │      │  Matching Engine │ ... │  Matching Engine │
           │  shard BTC/USDT  │      │  shard ETH/USDT  │     │  shard AAPL      │
           │  (in-memory      │      │                  │     │                  │
           │   order book)    │      │                  │     │                  │
           └─────┬────────┬───┘      └────────┬─────────┘     └────────┬─────────┘
                 │        │                   │                        │
       append-only│   emit │fills           ...                       ...
                 ▼        ▼
          ┌──────────────┐ ┌──────────────────┐   ┌─────────────────────┐
          │  Journal     │ │  Market Data     │   │  Trade Capture      │
          │  (NVMe +     │ │  Publisher       │──>│  (to settlement,    │
          │   raft x3)   │ │  (multicast/WS)  │   │   clearing, DWH)    │
          └──────────────┘ └──────────────────┘   └─────────────────────┘

Ключевой принцип: один поток (goroutine) на символ. Нет блокировок, нет гонок, детерминированный порядок. Масштабирование -- через шардирование по символу.

API дизайн

FIX (краткий пример)

FIX -- SOH-делимитированный текстовый протокол, стандарт для institutional trading.

8=FIX.4.4|9=176|35=D|49=CLIENT12|56=EXCHANGE|34=215|52=20260418-10:15:42.123|
11=order-42|21=1|55=BTCUSDT|54=1|60=20260418-10:15:42.100|38=0.5|40=2|44=68000.50|
59=0|10=128|

Поля: 35=D -- New Order Single, 54=1 -- Buy, 38 -- qty, 40=2 -- Limit, 44 -- price, 59=0 -- Day.

Для retail применяют WebSocket/JSON или бинарные протоколы (ITCH/OUCH, SBE).

REST (control plane)

POST /api/v1/orders
GET  /api/v1/orders/{order_id}
DELETE /api/v1/orders/{order_id}
GET  /api/v1/book/{symbol}?depth=20

DTO

<?php

declare(strict_types=1);

enum OrderSide: int
{
    case Buy = 1;
    case Sell = 2;
}

enum OrderType: int
{
    case Limit     = 1;
    case Market    = 2;
    case StopLimit = 3;
    case IOC       = 4;
    case FOK       = 5;
}

final readonly class NewOrderRequest
{
    public function __construct(
        public string $clientOrderId,
        public string $symbol,
        public OrderSide $side,
        public OrderType $type,
        public int $price,    // in ticks, integer math only
        public int $quantity,
        public string $timeInForce,
    ) {}
}

final readonly class ExecutionReport
{
    public function __construct(
        public int $orderId,
        public string $clientOrderId,
        public string $status,
        public int $filledQty,
        public int $avgPrice,
        public \DateTimeImmutable $ts,
    ) {}
}
Важно: **все цены и объёмы -- integers в тиках**. Float приводит к ошибкам округления, а на бирже одна копейка -- это юридическая проблема.

Схема данных

Matching engine держит состояние в памяти. Persistence -- только journal + периодические snapshots.

-- Journal (append-only, WAL-подобный). Хранится на NVMe с fsync.
CREATE TABLE journal (
    seq           BIGINT      PRIMARY KEY,   -- global sequence per symbol
    symbol        TEXT        NOT NULL,
    event_type    SMALLINT    NOT NULL,      -- 1=new, 2=cancel, 3=trade, 4=modify
    payload       BYTEA       NOT NULL,      -- binary-encoded event
    ts            TIMESTAMPTZ NOT NULL
);
CREATE INDEX ON journal (symbol, seq);

-- Snapshots (для быстрого recovery, не проигрывать весь journal).
CREATE TABLE snapshots (
    symbol        TEXT        NOT NULL,
    last_seq      BIGINT      NOT NULL,
    state         BYTEA       NOT NULL,      -- сериализованный order book
    taken_at      TIMESTAMPTZ NOT NULL,
    PRIMARY KEY (symbol, last_seq)
);

-- Trade capture (для клиринга, аналитики, налоговой отчётности).
CREATE TABLE trades (
    id            BIGINT      PRIMARY KEY,
    symbol        TEXT        NOT NULL,
    price         BIGINT      NOT NULL,
    qty           BIGINT      NOT NULL,
    buy_order_id  BIGINT      NOT NULL,
    sell_order_id BIGINT      NOT NULL,
    executed_at   TIMESTAMPTZ NOT NULL
) PARTITION BY RANGE (executed_at);

Journal пишется синхронно на 3 реплики (Raft-like). Пока кворум не подтвердил -- gateway не шлёт ack клиенту.

Ключевые компоненты

1. Order Book

Структура стакана должна поддерживать:

  • O(log N) вставку нового уровня цены;
  • O(1) доступ к best bid / best ask;
  • O(1) matching против топа;
  • O(qty) при частичном исполнении (FIFO по времени внутри уровня).

Классический вариант -- skip list или RB-tree по цене, плюс FIFO-очередь заявок на каждом уровне. Для узкого price range популярен bucket-based array (цена → индекс).

package engine

import (
    "container/list"
    "sync/atomic"
)

// Order represents a single resting order at a price level.
type Order struct {
    ID        uint64
    ClientID  string
    Side      OrderSide
    Price     int64
    Qty       int64
    Remaining int64
    SeqNo     uint64 // for strict time priority (cheaper than time.Time)
}

// priceLevel holds a FIFO queue of orders at the same price.
type priceLevel struct {
    Price  int64
    Orders *list.List // list.Element.Value = *Order
    Total  int64      // aggregate remaining qty for fast L2 snapshots
}

// OrderBook stores bids and asks for a single symbol. Not thread-safe:
// each symbol runs on its own goroutine.
type OrderBook struct {
    Symbol string
    Bids   *SkipList // sorted DESC by price
    Asks   *SkipList // sorted ASC by price
    orders map[uint64]*Order
    nextID atomic.Uint64
}

// Match attempts to fill an incoming order against the opposite side.
// Returns trades executed and the remaining order (if any).
func (ob *OrderBook) Match(in *Order) ([]Trade, *Order) {
    trades := make([]Trade, 0, 4)
    book := ob.Asks
    if in.Side == SideSell {
        book = ob.Bids
    }

    for in.Remaining > 0 {
        top := book.Top()
        if top == nil {
            break
        }
        // Price check depending on side
        if in.Side == SideBuy && top.Price > in.Price {
            break
        }
        if in.Side == SideSell && top.Price < in.Price {
            break
        }

        for e := top.Orders.Front(); e != nil && in.Remaining > 0; {
            maker := e.Value.(*Order)
            fillQty := min64(in.Remaining, maker.Remaining)

            trades = append(trades, Trade{
                Price:    maker.Price,
                Qty:      fillQty,
                TakerID:  in.ID,
                MakerID:  maker.ID,
            })
            in.Remaining -= fillQty
            maker.Remaining -= fillQty
            top.Total -= fillQty

            next := e.Next()
            if maker.Remaining == 0 {
                top.Orders.Remove(e)
                delete(ob.orders, maker.ID)
            }
            e = next
        }

        if top.Orders.Len() == 0 {
            book.RemoveTop()
        }
    }

    if in.Remaining == 0 {
        return trades, nil
    }
    return trades, in
}

// Add places a resting order into the book (after matching did not fully fill).
func (ob *OrderBook) Add(o *Order) {
    book := ob.Bids
    if o.Side == SideSell {
        book = ob.Asks
    }
    lvl := book.GetOrInsert(o.Price)
    lvl.Orders.PushBack(o)
    lvl.Total += o.Remaining
    ob.orders[o.ID] = o
}

func min64(a, b int64) int64 {
    if a < b {
        return a
    }
    return b
}
Skip list выбран потому, что:
  • вставка/удаление O(log N) без блокировок;
  • не ломает кэш процессора (в отличие от RB-tree с частыми ротациями);
  • легко пройти топ-N уровней для L2 snapshot.

Альтернатива для ликвидных символов с узким спредом -- массив price → *priceLevel, где индекс = (price - basePrice) / tick. Доступ O(1), но память = диапазон цен.

2. Goroutine per symbol

// SymbolShard is a single goroutine that owns the order book and a channel.
type SymbolShard struct {
    Symbol string
    In     chan Event  // serialized input
    Book   *OrderBook
    Out    chan<- Event
    Journal Journal
}

func (s *SymbolShard) Run() {
    for ev := range s.In {
        // 1. Journal first (sync, quorum write)
        if err := s.Journal.Append(ev); err != nil {
            s.Out <- Reject(ev, err)
            continue
        }
        // 2. Apply to book
        switch e := ev.(type) {
        case *NewOrderEvent:
            trades, rest := s.Book.Match(e.Order)
            for _, t := range trades {
                s.Out <- &TradeEvent{Trade: t}
            }
            if rest != nil {
                s.Book.Add(rest)
            }
        case *CancelEvent:
            s.Book.Cancel(e.OrderID)
        }
    }
}

Никаких мьютексов, никакого shared state с другими шардами. Latency -- только: journal write + memory access.

3. Sequencer

Чтобы гарантировать справедливый порядок, перед matching engine стоит sequencer -- одиночный компонент, который каждой заявке присваивает монотонный seq. Все gateway шлют заявки в sequencer; он пишет их в journal и раскидывает по шардам. Sequencer -- single point of truth по порядку, но он очень простой и легко реплицируется через Raft.

4. Recovery

При старте шарда:

  1. Загрузить последний snapshot (state + last_seq).
  2. Прочитать journal начиная с last_seq + 1 и applied events.
  3. Готов обслуживать новые.

Snapshot каждые N секунд / M событий. Это компромисс между размером journal tail и I/O.

5. Gateway и risk check

Gateway вне hot path. Он:

  • проверяет аутентификацию,
  • балансирует, достаточно ли залога / лимитов (pre-trade risk),
  • валидирует поля (tick size, lot size),
  • форвардит в sequencer.

Pre-trade risk -- отдельный сервис с in-memory кэшем позиций клиента. Асинхронные апдейты позиций после fills.

Масштабирование и bottlenecks

Bottleneck 1: один символ -- один поток

Если BTC/USDT генерирует 500K orders/sec, но один поток не тянет -- всё. Это фундаментальное ограничение price-time priority: нельзя параллелить matching без потери детерминизма.

Митигации:

  • оптимизация кода (избегать аллокаций в hot path, preallocated pools, jump tables вместо switch);
  • lock-free ring buffer вместо chan (Disruptor pattern, LMAX);
  • pin goroutine на CPU core, huge pages, отключить GC в hot path (debug.SetGCPercent(-1) и ручное управление);
  • использовать FPGA / kernel bypass (DPDK, Solarflare) для сетевого уровня.

LMAX Disruptor реально даёт 6M TPS на одном потоке -- с учётом pre-allocated events и cache-line padding.

Bottleneck 2: journal I/O

100K events/sec * 128 байт = 12.5 MB/s. Это немного, но синхронная репликация на 3 узла через сеть -- это latency 100-300 us только на fsync+network roundtrip.

Митигации:

  • group commit: батчить несколько заявок в один fsync (с осторожностью для p99);
  • NVMe + RDMA между репликами;
  • separate journal disk от snapshots.

Bottleneck 3: market data fan-out

100K подписчиков * 100 updates/sec на каждого = 10M messages/sec. Juniper-level.

Митигации:

  • UDP multicast в closed LAN для профессионалов;
  • conflated feed для retail (раз в 100ms);
  • edge CDN для WebSocket (но WS = order);
  • бинарные протоколы (SBE) вместо JSON.

Горизонтальное масштабирование

  • Шардирование по символу: раскидываем тысячи символов по десяткам шардов.
  • Active-passive replication для каждого шарда: primary принимает, 2 follower'а reply журнал. Failover < 1s через shared journal.
  • Нельзя шардировать один символ -- это сломает price-time priority.

Trade-offs и альтернативы

Решение Плюс Минус
In-memory book Latency < 10us Recovery стоит snapshots + journal
Skip list O(log N), кэш-дружественный Больше памяти чем RB-tree
Bucket array price levels O(1) Только для узкого диапазона
Goroutine per symbol Детерминизм, нет блокировок Нельзя шардировать один символ
FIX over TCP Индустриальный стандарт Verbose, нужно парсить SOH
SBE/ITCH Компактно, предсказуемо Свой формат на каждую биржу
Synchronous replication Zero loss +latency
Async replication Низкая latency Риск loss при split brain

Альтернативный матчер -- pro-rata / hybrid. Вместо FIFO внутри уровня раздаём fill пропорционально объёмам. Используется в фьючерсных контрактах. Менее справедлив для маленьких заявок, но ускоряет large blocks.

Auction-based matching (open/close auction) -- отдельный режим: собираем заявки период, вычисляем clearing price, фиксируем. Другой код path.

Выводы

Matching engine -- это редкий случай в enterprise-разработке, где счёт идёт на микросекунды и где шардировать нужно аккуратно, чтобы не сломать справедливость. Правильная архитектура -- один поток на символ + внешний sequencer + sync journal. В остальном применимы стандартные приёмы: избавляться от аллокаций, использовать integers, писать в immutable log. Всё сложное -- вокруг engine: pre-trade risk, market data fan-out, clearing, surveillance. Сам движок должен быть простым, детерминированным и предсказуемым.