Matching engine -- сердце любой биржи. Именно он сопоставляет заявки покупателей и продавцов, формирует сделки и поддерживает стакан (order book). От его корректности и предсказуемой latency зависит справедливость торгов и доверие участников рынка.
Функциональные требования
- Приём заявок (orders) разных типов: limit, market, stop-limit, IOC (immediate-or-cancel), FOK (fill-or-kill), post-only.
- Отмена и модификация заявок.
- Сопоставление заявок по принципу price-time priority (сначала лучшая цена, при равных ценах -- кто раньше поставил).
- Публикация market data: топ-уровня (L1), полного стакана (L2), trade tape (исполненные сделки).
- Журналирование всех событий (orders, trades, cancels) в append-only лог -- для аудита и recovery.
- Поддержка нескольких инструментов (symbols): BTC/USDT, ETH/USDT, AAPL, и т.д.
Нефункциональные требования
- Latency: p99 end-to-end tick-to-trade < 1 ms, медиана < 100 us внутри matching engine.
- Throughput: 100K заявок/сек на один инструмент, 1M+ на кластер.
- Fairness: детерминированный порядок обработки, никакой параллельности внутри символа.
- Durability: zero data loss для принятых заявок (synchronous replication journal).
- Availability: 99.995% в торговые часы, failover < 1s.
- 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,
) {}
}
Схема данных
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
}
- вставка/удаление 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
При старте шарда:
- Загрузить последний snapshot (
state + last_seq). - Прочитать journal начиная с
last_seq + 1и applied events. - Готов обслуживать новые.
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. Сам движок должен быть простым, детерминированным и предсказуемым.