Розробка системи збору order book у реальному часі

Розробка системи збору order book у реальному часі Уявіть: ваш торговий бот здійснює угоду за ціною, яка вже змінилася — книга заявок розсинхронізовано через обрив WebSocket. Втрати на одній такій події можуть сягати 2-3% від капіталу. Ми стикалися з цим на ранніх проектах і розробили систему, що

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

Часті запитання

Останні роботи

  • image_website-b2b-advance_0.webp
    Розробка сайту компанії B2B ADVANCE
    1451
  • image_web-applications_feedme_466_0.webp
    Розробка веб-додатків для компанії FEEDME
    1309
  • 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
    1011

Розробка системи збору order book у реальному часі

Уявіть: ваш торговий бот здійснює угоду за ціною, яка вже змінилася — книга заявок розсинхронізовано через обрив WebSocket. Втрати на одній такій події можуть сягати 2-3% від капіталу. Ми стикалися з цим на ранніх проектах і розробили систему, що виключає подібні інциденти.

Дані order book критичні для трьох сценаріїв: побудова торгового бота, створення агрегатора ліквідності, моніторинг ринку. У кожному з них спільна проблема — стабільне отримання high-frequency даних без втрат і з мінімальною латентністю. Невірний вибір протоколу або відсутність обробки reconnect’ів веде до розсинхронізації стакану та збиткових угод. За 5+ років роботи ми реалізували понад 30 інтеграцій з різними біржами. Гарантуємо стабільність збору даних та низьку затримку.

Чому WebSocket краще за REST?

REST polling (GET /api/v3/depth?symbol=BTCUSDT) — неправильний вибір для real-time order book. На активних ринках книга заявок оновлюється 10–100 разів на секунду. Поллінг раз на секунду дає застарілі дані та навантажує API ліміти. Правильний підхід — WebSocket потоки з інкрементальними оновленнями, які у 10-100 разів швидші.

Покроковий алгоритм:

  1. Запит snapshot через REST (повний стакан на поточний момент)
  2. Підписка на WebSocket потік оновлень
  3. Застосування інкрементальних оновлень до snapshot
  4. Перевірка sequence ID – при виявленні gap виконується ресинхронізація
Параметр REST Polling WebSocket Stream
Затримка 1 сек+ (період опитування) 10-100 мс (подійно)
Навантаження на API Висока (запити щосекунди) Низька (одне підключення)
Актуальність даних Миттєво застаріває Завжди останній стан
Масштабування Проблеми при кількох інструментах До 1024 потоків на ключ

Більшість великих CEX (Binance, Bybit, OKX) дотримуються однієї схеми:

  1. Отримати snapshot через REST (повний стакан на поточний момент)
  2. Підписатися на WebSocket потік оновлень
  3. Застосовувати оновлення до snapshot, підтримуючи локальну копію стакану

Деталі протоколу: Binance Diff. Depth Stream

import asyncio, json, aiohttp from sortedcontainers import SortedDict class OrderBook: def __init__(self): self.bids = SortedDict(lambda x: -x) # спадаючий порядок self.asks = SortedDict() self.last_update_id = 0 def apply_update(self, bids: list, asks: list, update_id: int): if update_id <= self.last_update_id: return # застаріле оновлення, ігноруємо for price, qty in bids: price, qty = float(price), float(qty) if qty == 0: self.bids.pop(price, None) # видалити рівень else: self.bids[price] = qty for price, qty in asks: price, qty = float(price), float(qty) if qty == 0: self.asks.pop(price, None) else: self.asks[price] = qty self.last_update_id = update_id @property def best_bid(self) -> tuple[float, float] | None: if self.bids: price = self.bids.keys()[0] return price, self.bids[price] return None @property def best_ask(self) -> tuple[float, float] | None: if self.asks: price = self.asks.keys()[0] return price, self.asks[price] return None 

Як синхронізувати стакан після обриву?

При обриві з'єднання або втраті пакетів ризик неконсистентності високий. Ми використовуємо наступні техніки:

  • Буферизація оновлень до отримання знімка (як показано вище)
  • Перевірка sequence ID кожного оновлення: якщо update_id не збігається з очікуваним, відкидаємо пакет і запитуємо новий знімок
  • Експоненційний backoff при перепідключенні з обмеженням у 60 секунд
  • Моніторинг затримки з оповіщенням при перевищенні порогу (наприклад, >500 мс)
async def connect_ws_with_retry(url: str, handler, max_retries=10): for attempt in range(max_retries): try: async with websockets.connect(url, ping_interval=20) as ws: async for message in ws: await handler(message) except (websockets.exceptions.ConnectionClosed, Exception) as e: wait = min(2 ** attempt, 60) # max 60 секунд logging.warning(f"WS disconnected: {e}, retry in {wait}s") await asyncio.sleep(wait) 
Деталі роботи з Binance depth stream

Binance — найчастіший запит. У них два stream варіанти:

  • btcusdt@depth — оновлення кожні 100ms або 1000ms (параметр @depth@100ms)
  • btcusdt@depth20 — топ-20 рівнів кожні 100ms (без інкрементальних оновлень, завжди повний)

Для повного стакану з застосуванням патчів:

async def maintain_binance_orderbook(symbol: str): ob = OrderBook() buffer = [] # буфер оновлень до отримання snapshot async def handle_ws_message(msg): data = json.loads(msg) # Накопичуємо оновлення ПОКИ не отримаємо snapshot if ob.last_update_id == 0: buffer.append(data) return # Binance: оновлення валідне якщо U <= lastUpdateId+1 <= u if data['U'] <= ob.last_update_id + 1 <= data['u']: ob.apply_update(data['b'], data['a'], data['u']) # Запускаємо WS ws_task = asyncio.create_task(connect_ws( f"wss://stream.binance.com:9443/ws/{symbol.lower()}@depth@100ms", handle_ws_message )) # Отримуємо snapshot (трохи чекаємо щоб буфер накопичився) await asyncio.sleep(0.5) async with aiohttp.ClientSession() as session: async with session.get( f"https://api.binance.com/api/v3/depth", params={"symbol": symbol.upper(), "limit": 1000} ) as resp: snapshot = await resp.json() # Ініціалізуємо стакан із snapshot for price, qty in snapshot['bids']: ob.bids[float(price)] = float(qty) for price, qty in snapshot['asks']: ob.asks[float(price)] = float(qty) ob.last_update_id = snapshot['lastUpdateId'] # Застосовуємо буферизовані оновлення for update in buffer: if update['u'] > ob.last_update_id: ob.apply_update(update['b'], update['a'], update['u']) await ws_task 

Критичний момент: якщо пропущено update (gap в U → u послідовності) — стакан розсинхронізовано. Потрібна логіка ресинхронізації: детектувати gap та переініціалізувати з нового snapshot.

Кейс: агрегація Binance та Bybit для арбітражу

Для крос-біржового арбітражу необхідно підтримувати книги заявок кількох бірж паралельно. Ось приклад агрегатора, який знаходить найкращу ціну:

EXCHANGES = { "binance": BinanceOrderBook, "bybit": BybitOrderBook, "okx": OKXOrderBook, } async def run_aggregator(symbol: str): books = {name: cls(symbol) for name, cls in EXCHANGES.items()} tasks = [book.run() for book in books.values()] await asyncio.gather(*tasks) def get_best_price_across_exchanges(books: dict[str, OrderBook]) -> dict: best_bids = [(name, *ob.best_bid) for name, ob in books.items() if ob.best_bid] best_asks = [(name, *ob.best_ask) for name, ob in books.items() if ob.best_ask] best_bids.sort(key=lambda x: x[1], reverse=True) best_asks.sort(key=lambda x: x[1]) return { "best_bid": {"exchange": best_bids[0][0], "price": best_bids[0][1], "qty": best_bids[0][2]}, "best_ask": {"exchange": best_asks[0][0], "price": best_asks[0][1], "qty": best_asks[0][2]}, "spread": best_asks[0][1] - best_bids[0][1] } 

Агрегація в реальному часі дозволяє бачити найкращу ціну на купівлю та продаж по всіх біржах. Це основа для арбітражних стратегій та побудови unified order book.

Зберігання даних: TimescaleDB vs файлове

Для backtesting або аудиту — зберігання потоку оновлень, а не тільки знімків. L2 order book updates — це великий обсяг: для BTC/USDT на Binance ~100MB/год нестиснутих даних (для 10 пар ~1GB/год).

Критерій TimescaleDB Файлове (Parquet)
Запити в реальному часі Так (SQL) Ні (тільки аналітика)
Стиснення Автоматичне Налаштовуване (lz4)
Відтворення потоку Потрібна додаткова обробка Пряме читання
Інтеграція з Kafka Є Відсутня

Ми рекомендуємо використовувати TimescaleDB для довгострокового зберігання та Parquet для аналітики. За необхідності інтегруємо потік в Kafka для downstream-систем.

# Запис у бінарний формат через msgpack import msgpack, lz4.frame def serialize_update(update: dict) -> bytes: packed = msgpack.packb(update, use_bin_type=True) return lz4.frame.compress(packed) # TimescaleDB для time-series зберігання # Гіпертаблиця автоматично партиціонує за часом CREATE TABLE ob_updates ( time TIMESTAMPTZ NOT NULL, exchange TEXT NOT NULL, symbol TEXT NOT NULL, side CHAR(1) NOT NULL, -- 'b' або 'a' price NUMERIC NOT NULL, quantity NUMERIC NOT NULL ); SELECT create_hypertable('ob_updates', 'time'); 

Що входить у роботу?

При замовленні системи збору order book ви отримуєте:

  • Архітектурне рішення з вибором протоколів та стратегії ресинхронізації
  • Вихідний код на Python з asyncio та документацію по розгортанню
  • Дашборди Grafana для моніторингу затримок та помилок
  • Гарантію стабільності 99.9% та підтримку протягом двох тижнів після впровадження

Розробка системи під ключ для однієї біржі стартує від $2000, для трьох — від $5000. Термін виконання — 2-4 тижні. Оцінка проекту безкоштовно — зв'яжіться з нами для обговорення.

Типові помилки при зборі order book

  • Ігнорування sequence ID та відсутність ресинхронізації після gap
  • Використання REST polling замість WebSocket (призводить до затримок та лімітів)
  • Неправильний порядок: спочатку snapshot, потім підписка на оновлення
  • Відсутність буферизації оновлень до отримання snapshot (втрачаються перші пакети)
  • Некоректна обробка reconnect з експоненційним backoff

Замовте розробку системи збору order book з гарантією стабільності та низькою латентністю. Наш досвід у high-frequency trading дозволяє створювати рішення, які не втрачають дані та не розсинхронізуються навіть при пікових навантаженнях.