Реализация real-time WebSocket scraping для криптобирж и EVM-сетей

Представьте: вы торгуете на Binance, используя REST polling раз в секунду. За это время цена могла измениться на 0.5%, а вы пропустили арбитраж. WebSocket-подписка даёт событие в момент его возникновения — задержка снижается с 500 мс до 10–50 мс. Для мониторинга цен, order book и on-chain событий ра

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

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

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

  • 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

Представьте: вы торгуете на Binance, используя REST polling раз в секунду. За это время цена могла измениться на 0.5%, а вы пропустили арбитраж. WebSocket-подписка даёт событие в момент его возникновения — задержка снижается с 500 мс до 10–50 мс. Для мониторинга цен, order book и on-chain событий разница принципиальна.

Polling REST API раз в N секунд — неправильный инструмент для задач, требующих реакции на события. При polling с интервалом 1 секунда средняя задержка обнаружения события — 0.5 секунды. WebSocket-подписка даёт событие в момент его возникновения, задержка определяется только сетью (10–50 мс до ближайшего сервера биржи). Для мониторинга цен, order book и on-chain событий разница принципиальна. Один из наших клиентов сократил latency с 800 мс до 30 мс, внедрив WebSocket scraping для 50 пар на 5 биржах — экономия составила до 40% упущенной прибыли.

Параметр REST Polling WebSocket
Задержка события 500 мс – 2 с 10–50 мс
Нагрузка на сервер Высокая (N запросов/мин) Низкая (одно соединение)
Реакция на изменения С задержкой, возможны пропуски Мгновенная, все события подряд
Сложность реализации Низкая Средняя, требуется reconnect logic

Почему WebSocket scraping лучше REST polling для real-time данных?

WebSocket сокращает задержку в 10 раз по сравнению с REST polling (с 0.5 с до 50 мс) — это экономия до 40% упущенной прибыли. Для криптотрейдинга и DeFi-ботов эта разница критична.

Как настроить WebSocket-соединение с биржей?

Каждая биржа имеет свой протокол подписки. Паттерны схожи, детали различаются.

Binance: stream names через symbol@streamType

import asyncio import json import websockets async def binance_stream(symbols: list[str]): streams = '/'.join([f"{s.lower()}@trade" for s in symbols]) url = f"wss://stream.binance.com:9443/stream?streams={streams}" async with websockets.connect(url, ping_interval=20, ping_timeout=10) as ws: async for message in ws: data = json.loads(message) stream_data = data.get('data', data) yield { 'exchange': 'binance', 'symbol': stream_data['s'], 'price': float(stream_data['p']), 'amount': float(stream_data['q']), 'timestamp': stream_data['T'], 'is_buyer_maker': stream_data['m'], } 

Coinbase Advanced Trade: subscribe с channel и product_ids

subscribe_msg = { "type": "subscribe", "channel": "ticker", "product_ids": ["BTC-USD", "ETH-USD"], } 

Kraken

Использует генерацию subscription ID и имеет особенности формата ответа с парой в массиве. Детали описаны в официальной документации Kraken WebSocket API.

Ethereum/EVM: WebSocket подписки через web3.py

On-chain события через WebSocket subscriptions к Ethereum ноде (Alchemy, Infura, QuickNode или собственная нода):

from web3 import AsyncWeb3, WebSocketProvider async def subscribe_to_transfers(token_address: str): w3 = AsyncWeb3(WebSocketProvider( "wss://eth-mainnet.g.alchemy.com/v2/YOUR_KEY" )) # ERC-20 Transfer event signature hash transfer_sig = w3.keccak(text="Transfer(address,address,uint256)").hex() subscription_id = await w3.eth.subscribe('logs', { 'address': token_address, 'topics': [transfer_sig] }) async for payload in w3.socket.process_subscriptions(): if payload['subscription'] == subscription_id: log = payload['result'] yield decode_transfer_log(log) 

Ethereum JSON-RPC WebSocket поддерживает три типа подписок: newHeads (новые блоки), logs (события контрактов), newPendingTransactions (mempool транзакции). Подробнее в официальной документации Ethereum.

Почему важен reconnect и staleness watchdog?

WebSocket соединения разрываются по разным причинам: timeout сервера, network hiccup, перезапуск сервиса биржи. Production система должна автоматически восстанавливаться:

import asyncio import websockets from datetime import datetime class RobustWebSocketClient: def __init__(self, url: str, reconnect_delay: float = 1.0): self.url = url self.reconnect_delay = reconnect_delay self.max_reconnect_delay = 60.0 self.last_message_at = None self.stale_threshold = 30 # секунд без сообщений = staleness async def connect_with_retry(self, on_message, on_subscribe): delay = self.reconnect_delay while True: try: async with websockets.connect( self.url, ping_interval=20, ping_timeout=10, close_timeout=5, ) as ws: await on_subscribe(ws) delay = self.reconnect_delay # сбрасываем при успехе async for msg in ws: self.last_message_at = datetime.utcnow() await on_message(msg) except (websockets.ConnectionClosed, websockets.InvalidHandshake, OSError) as e: print(f"Connection error: {e}, reconnecting in {delay}s") await asyncio.sleep(delay) delay = min(delay * 2, self.max_reconnect_delay) async def staleness_watchdog(self): """Детектирует зависшее соединение без явного разрыва""" while True: await asyncio.sleep(10) if self.last_message_at: elapsed = (datetime.utcnow() - self.last_message_at).seconds if elapsed > self.stale_threshold: raise RuntimeError(f"Connection stale: {elapsed}s without data") 

Экспоненциальная задержка переподключения и watchdog на stale connection — обязательный минимум для промышленного scraping.

Как управлять order book через WebSocket?

Большинство бирж отдают order book через incremental updates — только изменившиеся уровни. Локальное поддержание актуального состояния order book:

from sortedcontainers import SortedDict class LocalOrderBook: def __init__(self): self.bids = SortedDict(lambda k: -k) # descending self.asks = SortedDict() # ascending self.last_update_id = 0 def apply_snapshot(self, snapshot: dict): self.bids.clear() self.asks.clear() for price, qty in snapshot['bids']: self.bids[float(price)] = float(qty) for price, qty in snapshot['asks']: self.asks[float(price)] = float(qty) self.last_update_id = snapshot['lastUpdateId'] def apply_update(self, update: dict): if update['u'] <= self.last_update_id: return # устаревший update, игнорируем for price, qty in update['b']: # bids p, q = float(price), float(qty) if q == 0: self.bids.pop(p, None) else: self.bids[p] = q for price, qty in update['a']: # asks p, q = float(price), float(qty) if q == 0: self.asks.pop(p, None) else: self.asks[p] = q self.last_update_id = update['u'] def best_bid(self) -> tuple[float, float]: k = next(iter(self.bids)) return k, self.bids[k] def best_ask(self) -> tuple[float, float]: k = next(iter(self.asks)) return k, self.asks[k] 

Важно: при старте нужно получить снапшот через REST, затем применять WebSocket updates начиная с lastUpdateId > snapshotId. Обновления до снапшота отбрасываются, пропуск в последовательности Uu требует повторного снапшота.

Масштабирование: множество пар и бирж

Одна async event loop в Python справляется с 50–200 одновременными WebSocket соединениями. Для большего числа — несколько процессов или Go-сервис (goroutines значительно легче asyncio tasks).

Fanout результатов: обработанные сообщения публикуются в Redis Pub/Sub или Kafka для downstream consumers. WebSocket handler должен минимально обрабатывать данные и быстро публиковать — тяжёлую обработку делает отдельный consumer.

Мониторинг здоровья

Метрики для каждого WebSocket соединения: messages per second, reconnect count, last message timestamp, lag от биржевого timestamp до processing timestamp. Используем Grafana + Prometheus alerting на stale connections (> 60 сек без сообщений по активной паре).

Метрика Описание Порог алерта
messages/sec Количество сообщений в секунду < 0.5 ожидаемого
reconnects Количество переподключений за час > 5
last_message_age Время с последнего сообщения > 60 с
lag Задержка от биржевого времени > 500 мс

Что входит в настройку WebSocket scraping

  • Подключение к биржам / блокчейн-нодам по WebSocket (Binance, Coinbase, Kraken, Ethereum, Polygon, Solana и др.)
  • Реализация reconnect logic с exponential backoff и staleness watchdog
  • Локальная агрегация order book с синхронизацией через снапшоты
  • Публикация нормализованных данных в Redis Pub/Sub или Kafka
  • Мониторинг и алерты (Grafana, Prometheus)
  • Документация по архитектуре и настройке
  • Обучение вашей команды работе с системой

Наш опыт и гарантии

За 5+ лет мы реализовали 50+ проектов real-time scraping для криптобирж, DeFi-протоколов и NFT-маркетплейсов. Гарантируем стабильную работу, автоматическое восстановление после сбоев и мониторинг 24/7. Работаем с Ethereum, Binance, Polygon, Arbitrum, Solana и другими сетями.

Настройка real-time парсинга для 3–5 бирж с мониторингом 20–50 пар, reconnect логикой и публикацией в Redis/Kafka занимает 1–2 дня. Свяжитесь с нами для расчёта стоимости. Закажите настройку уже сегодня — получите консультацию.