Разработка системы граббинга order book в реальном времени

Парсинг данных order book с бирж в реальном времени Представьте: ваш торговый бот совершает сделку по цене, которая уже изменилась — стакан рассинхронизирован из-за обрыва WebSocket. Потери на одном таком событии могут достигать 2-3% от капитала. Мы сталкивались с этим на ранних проектах и разраб

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

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

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

  • image_website-b2b-advance_0.webp
    Разработка сайта компании B2B ADVANCE
    1452
  • 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’ов ведёт к рассинхронизации стакана и убыточным сделкам. За годы работы мы реализовали свыше 30 интеграций с различными биржами. Гарантируем стабильность сбора данных и низкую задержку.

Почему WebSocket лучше REST?

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

Параметр REST Polling WebSocket Stream
Задержка 1 сек+ (период опроса) 10-100 мс (событийно)
Нагрузка на API Высокая (запросы каждую секунду) Низкая (одно подключение)
Актуальность данных Мгновенно устаревает Всегда последнее состояние
Масштабирование Проблемы при нескольких инструментах До 1024 потоков на ключ

Большинство крупных CEX (Binance, Bybit, OKX) следуют одной схеме:

  1. Получить snapshot через REST (полный стакан на текущий момент)
  2. Подписаться на WebSocket поток обновлений
  3. Применять обновления к snapshot, поддерживая локальную копию стакана
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 в Uu последовательности) — стакан рассинхронизирован. Нужна логика ресинхронизации: детектировать 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/час несжатых данных.

Критерий 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% и поддержку в течение двух недель после внедрения

Сроки разработки: от 2 до 4 недель в зависимости от количества бирж и требуемой архитектуры. Стоимость рассчитывается индивидуально. Свяжитесь с нами для обсуждения вашего проекта — мы подберём оптимальное решение.

Типичные ошибки при граббинге order book

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

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