Разработка pipeline обработки order book данных для ML

Полный стакан ордеров содержит всю информацию о ликвидности на бирже. Но собрать его, нормализовать и превратить в признаки для машинного обучения — нетривиальная инженерная задача. Мы построили production-grade pipeline для Binance, Bybit и OKX, который обрабатывает до 10 000 обновлений в секунду.

Направления блокчейн-разработки

Часто задаваемые вопросы

Последние работы

  • image_website-b2b-advance_0.webp
    Разработка сайта компании B2B ADVANCE
    1452
  • image_web-applications_feedme_466_0.webp
    Разработка веб-приложения для компании FEEDME
    1310
  • image_websites_belfingroup_462_0.webp
    Разработка веб-сайта для компании БЕЛФИНГРУПП
    1005
  • image_ecommerce_furnoro_435_0.webp
    Разработка интернет магазина для компании FURNORO
    1270
  • image_logo-advance_0.webp
    Разработка логотипа компании B2B Advance
    719
  • image_crm_enviok_479_0.webp
    Разработка веб-приложения для компании Enviok
    1012

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

  1. Установка соединения: через wss://stream.binance.com:9443/ws/btcusdt@depth@100ms (аналог для других бирж).
  2. Первоначальный REST-снимок: синхронизация через updateId для обеспечения консистентности.
  3. Инкрементальные обновления: каждое сообщение diff stream накладывается на текущее состояние стакана.
  4. Сохранение снэпшотов: с заданной периодичностью (каждое 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 рабочих дней. Точную оценку даём после бесплатного аудита ваших данных. Закажите анализ — и мы подберём оптимальную архитектуру под ваш объём стакана. Свяжитесь с нами для консультации.