Полный стакан ордеров содержит всю информацию о ликвидности на бирже. Но собрать его, нормализовать и превратить в признаки для машинного обучения — нетривиальная инженерная задача. Мы построили 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 рабочих дней. Точную оценку даём после бесплатного аудита ваших данных. Закажите анализ — и мы подберём оптимальную архитектуру под ваш объём стакана. Свяжитесь с нами для консультации.







