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







