Представьте: вы торгуете на 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. Обновления до снапшота отбрасываются, пропуск в последовательности U → u требует повторного снапшота.
Масштабирование: множество пар и бирж
Одна 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 дня. Свяжитесь с нами для расчёта стоимости. Закажите настройку уже сегодня — получите консультацию.







