Розробка pipeline обробки order book даних для ML
Повний стакан ордерів містить всю інформацію про ліквідність на біржі. Але зібрати його, нормалізувати і перетворити на ознаки для машинного навчання — нетривіальне інженерне завдання. Ми побудували production-grade pipeline для Binance, Bybit та OKX, який обробляє до 10 000 оновлень за секунду. Наш досвід включає інтеграцію з 15+ криптобіржами та зберігання близько 5 ТБ даних на місяць. Ми маємо 5+ років досвіду в криптоінфраструктурі та 50+ реалізованих проектів під ключ. Повний 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 на рік. Вартість розробки базового pipeline під ключ — від $5,000 до $15,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 робочих днів. Точну оцінку даємо після безкоштовного аудиту ваших даних. Замовте аналіз — і ми підберемо оптимальну архітектуру під ваш обсяг стакана. Пишіть нам для консультації — оцінимо проєкт безкоштовно.







