Полный стакан ордеров содержит всю информацию о ликвидности на бирже. Но собрать его, нормализовать и превратить в признаки для машинного обучения — нетривиальная инженерная задача. Мы построили production-grade pipeline для Binance, Bybit и OKX, который обрабатывает до 10 000 обновлений в секунду. Наш опыт включает интеграцию с 15+ криптобиржами и хранение порядка 5 ТБ данных в месяц. Полный L2 стакан описывает каждый уровень цены с объёмом — это основа для построения краткосрочных прогнозов. Гарантируем стабильный сбор при пиковых нагрузках и консистентность снэпшотов.
Заказчики часто приходят с сырыми WebSocket-стримами, не зная, как синхронизировать diff stream с REST-снимком. Ошибка в one-off приводит к разъезду стакана и неверным сигналам. Мы решаем эту проблему на уровне архитектуры коллектора.
Проблемы, которые решаем
- Объём данных. Полный L2 стакан на Binance содержит 5000 уровней с каждой стороны. При обновлениях каждые 100 мс это генерирует десятки гигабайт в день. Наивное хранение в PostgreSQL убьёт производительность.
- Гонка состояний. WebSocket diff stream приходит асинхронно. Без синхронизации с REST-снимком стакан разъезжается — цена уходит в несуществующие уровни.
- Формат данных. Каждая биржа отдаёт стакан по-своему: Binance — вложенные массивы, Coinbase — JSON с разными ключами. Нужен единый интерфейс.
Как синхронизировать WebSocket diff stream с REST-снимком?
Алгоритм прост: открываем WebSocket, получаем первый diff stream, сразу запрашиваем REST-снимок с полным состоянием. Далее каждое обновление накладываем на локальный стакан. Для контроля используем lastUpdateId: применяем только сообщения с u > lastUpdateId. Если последовательность нарушена — перезапрашиваем снимок. Этот подход исключает разъезд стакана даже при высокой волатильности.
Как собрать order book через WebSocket: пошаговый алгоритм
- Установка соединения: через
wss://stream.binance.com:9443/ws/btcusdt@depth@100ms(аналог для других бирж). - Первоначальный REST-снимок: синхронизация через
updateIdдля обеспечения консистентности. - Инкрементальные обновления: каждое сообщение diff stream накладывается на текущее состояние стакана.
- Сохранение снэпшотов: с заданной периодичностью (каждое N-е обновление) фиксируется полное состояние для последующего feature engineering.
Пример кода коллектора:
import asyncio
import websockets
import json
from collections import deque
class OrderBookCollector:
def __init__(self, symbol, max_depth=100):
self.symbol = symbol
self.bids = {}
self.asks = {}
self.max_depth = max_depth
self.snapshots = deque(maxlen=10000)
async def connect_binance(self):
url = f"wss://stream.binance.com:9443/ws/{self.symbol.lower()}@depth@100ms"
async with websockets.connect(url) as ws:
await self.fetch_snapshot()
async for msg in ws:
data = json.loads(msg)
self.process_diff_update(data)
if len(self.snapshots) % 10 == 0:
self.save_snapshot()
def process_diff_update(self, data):
for bid_level in data.get('b', []):
price, qty = float(bid_level[0]), float(bid_level[1])
if qty == 0:
self.bids.pop(price, None)
else:
self.bids[price] = qty
for ask_level in data.get('a', []):
price, qty = float(ask_level[0]), float(ask_level[1])
if qty == 0:
self.asks.pop(price, None)
else:
self.asks[price] = qty
def get_features(self, n_levels=20):
sorted_bids = sorted(self.bids.items(), reverse=True)[:n_levels]
sorted_asks = sorted(self.asks.items())[:n_levels]
if not sorted_bids or not sorted_asks:
return None
mid_price = (sorted_bids[0][0] + sorted_asks[0][0]) / 2
features = {}
for i, (price, qty) in enumerate(sorted_bids[:10]):
features[f'bid_qty_{i}'] = qty
features[f'bid_dist_{i}'] = (mid_price - price) / mid_price
for i, (price, qty) in enumerate(sorted_asks[:10]):
features[f'ask_qty_{i}'] = qty
features[f'ask_dist_{i}'] = (price - mid_price) / mid_price
bid_vol_n = sum(qty for _, qty in sorted_bids[:5])
ask_vol_n = sum(qty for _, qty in sorted_asks[:5])
features['obi_5'] = (bid_vol_n - ask_vol_n) / (bid_vol_n + ask_vol_n + 1e-8)
bid_vol_20 = sum(qty for _, qty in sorted_bids[:20])
ask_vol_20 = sum(qty for _, qty in sorted_asks[:20])
features['obi_20'] = (bid_vol_20 - ask_vol_20) / (bid_vol_20 + ask_vol_20 + 1e-8)
features['wmid'] = (sorted_bids[0][0] * sorted_asks[0][1] + sorted_asks[0][0] * sorted_bids[0][1]) / (sorted_bids[0][1] + sorted_asks[0][1])
features['spread'] = (sorted_asks[0][0] - sorted_bids[0][0]) / mid_price
for n in [5, 10, 20]:
bid_depth = sum(qty for _, qty in sorted_bids[:n])
ask_depth = sum(qty for _, qty in sorted_asks[:n])
features[f'depth_ratio_{n}'] = bid_depth / max(ask_depth, 1e-8)
return features
Почему ClickHouse — оптимальное хранилище для order book?
Полный L2 стакан — огромный объём. ClickHouse в 10 раз быстрее PostgreSQL на колоночных агрегациях. Согласно документации ClickHouse, колоночная СУБД обеспечивает сжатие до 10 раз и скорость записи более 1 млн строк в секунду. Сравните:
| СУБД | Скорость записи (строк/с) | Сжатие | Агрегации по времени |
|---|---|---|---|
| PostgreSQL | ~100 000 | 2-5x | Медленные |
| TimescaleDB | ~200 000 | 3-6x | Средние |
| ClickHouse | ~1 000 000 | 5-10x | Быстрые |
Пример схемы хранения с автоматическим TTL:
CREATE TABLE order_book_snapshots (
timestamp DateTime64(3),
symbol LowCardinality(String),
exchange LowCardinality(String),
bid_price_0 Float32, bid_qty_0 Float32,
bid_price_1 Float32, bid_qty_1 Float32,
-- ... до bid_price_19, bid_qty_19
ask_price_0 Float32, ask_qty_0 Float32,
-- ...
spread Float32,
obi_5 Float32,
obi_20 Float32
) ENGINE = MergeTree()
PARTITION BY toYYYYMMDD(timestamp)
ORDER BY (symbol, timestamp)
TTL timestamp + INTERVAL 90 DAY;
Экономия на инфраструктуре при использовании ClickHouse достигает 70% за счёт сжатия — это около $20,000 в год для проекта с 5 ТБ данных. Для крупных проектов экономия может достигать $30,000 в год.
Feature engineering из order book
На основе собранных снэпшотов строим признаки. Базовые: OBI (order book imbalance), спред, глубина. Дополнительные: скользящие средние OBI, его волатильность, кумулятивный поток ордеров (COF).
def engineer_orderbook_features(snapshots_df, window_sizes=[10, 50, 100]):
features = snapshots_df.copy()
for window in window_sizes:
features[f'obi_5_ma_{window}'] = features['obi_5'].rolling(window).mean()
features[f'obi_5_delta_{window}'] = features['obi_5'].diff(window)
features[f'obi_5_std_{window}'] = features['obi_5'].rolling(window).std()
features['cof'] = features['obi_5'].cumsum()
features['cof_ma'] = features['cof'].rolling(100).mean()
features['cof_deviation'] = features['cof'] - features['cof_ma']
features['spread_ma'] = features['spread'].rolling(50).mean()
features['spread_ratio'] = features['spread'] / features['spread_ma']
features['depth_change'] = features['depth_ratio_10'].diff(10)
return features
Как оценить качество прогноза mid-price?
Для краткосрочного прогноза mid-price (через N обновлений стакана) используем метрики accuracy, precision и F1-score для бинарной классификации направления движения. Код подготовки обучающей выборки:
def create_training_data(snapshots_df, prediction_horizon=10):
features = engineer_orderbook_features(snapshots_df)
future_mid = snapshots_df['mid_price'].shift(-prediction_horizon)
current_mid = snapshots_df['mid_price']
target = np.sign(future_mid - current_mid)
valid_mask = features.notna().all(axis=1) & target.notna()
return features[valid_mask], target[valid_mask]
Типичные ошибки при разработке order book pipeline
Даже опытные команды допускают ошибки: игнорирование перекосов стакана в моменты высокой волатильности, неправильная обработка событий lastUpdateId, отсутствие проверки консистентности после реконнекта. Мы сталкивались с проектом, где из-за пропущенных диффов стакан разошёлся на 20% — модель показывала ложные сигналы. Решение — встраивание проверок контрольных сумм и автоматическое восстановление полного снэпшота при обнаружении несоответствия.
Что входит в разработку pipeline
- Исходный код коллектора и пайплайна (асинхронный Python).
- Дампы тестовых данных для offline-тестирования.
- README с подробными примерами использования.
- Миграции схемы ClickHouse с TTL.
- Обучение вашей команды работе с pipeline.
Этапы работы и сроки
| Этап | Длительность | Результат |
|---|---|---|
| Аналитика | 2-3 дня | Спецификация API, объёмов |
| Проектирование | 2-3 дня | Схема хранения, выбор признаков |
| Реализация | 5-10 дней | Коллектор, пайплайн, код |
| Тестирование | 3-5 дней | Симуляция 24h, отчёты |
| Деплой | 2-3 дня | Docker, мониторинг |
Базовый pipeline для одной биржи с моделью LightGBM — от 14 до 30 рабочих дней. Точную оценку даём после бесплатного аудита ваших данных. Закажите анализ — и мы подберём оптимальную архитектуру под ваш объём стакана. Свяжитесь с нами для консультации.







