Уявіть: ви торгуєте на Binance, використовуючи REST polling раз на секунду. Ми допомагаємо клієнтам налаштовувати WebSocket scraping, щоб уникнути таких втрат. За 5+ років ми реалізували понад 50 проєктів з real-time scraping для провідних криптобірж та DeFi-протоколів.
Polling REST API раз на N секунд — неправильний інструмент для задач, що потребують реакції на події. При polling з інтервалом 1 секунда середня затримка виявлення події — 0.5 секунди. WebSocket-підписка дає подію в момент її виникнення, затримка визначається тільки мережею (10–50 мс до найближчого сервера біржі). WebSocket забезпечує в 10 разів меншу затримку, ніж REST polling. Для моніторингу цін, order book та on-chain подій різниця принципова. Один із наших клієнтів скоротив latency з 800 мс до 30 мс, впровадивши WebSocket scraping для 50 пар на 5 біржах — економія склала до $2000 щомісяця на комісіях.
| Параметр | REST Polling | WebSocket |
|---|---|---|
| Затримка події | 500 мс – 2 с | 10–50 мс |
| Навантаження на сервер | Висока (N запитів/хв) | Низька (одне з'єднання) |
| Реакція на зміни | Із затримкою, можливі пропуски | Миттєва, всі події підряд |
| Складність реалізації | Низька | Середня, потребує reconnect logic |
Чому WebSocket scraping кращий за REST polling для real-time даних?
WebSocket скорочує затримку в 10 разів порівняно з REST polling — це економія до 40% втраченого прибутку. Для криптотрейдингу та DeFi-ботів ця різниця критична.
Як налаштувати WebSocket-з'єднання з біржею?
Кожна біржа має свій протокол підписки. Патерни схожі, деталі відрізняються.
- Виберіть біржу та тип даних (наприклад, trade, ticker, order book).
- Знайдіть документацію WebSocket API на сайті біржі.
- Напишіть код підписки, використовуючи приклади нижче.
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.
Підтримуємо L2 rollup мережі, такі як Arbitrum та Optimism, через їхні WebSocket endpoints.
Чому важливі reconnect та staleness watchdog?
WebSocket з'єднання розриваються з різних причин: timeout сервера, network hiccup, перезапуск сервісу біржі. Production система повинна автоматично відновлюватися:
Приклад реалізації RobustWebSocketClient
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:
Приклад класу LocalOrderBook
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)
- Документація з архітектури та налаштування
- Навчання вашої команди роботі з системою
Наш досвід та гарантії
Маємо досвід реалізації понад 50 проєктів real-time scraping для криптобірж, DeFi-протоколів та NFT-маркетплейсів. Гарантуємо стабільну роботу, автоматичне відновлення після збоїв та моніторинг 24/7. Працюємо з Ethereum, Binance, Polygon, Arbitrum, Solana та іншими мережами.
Налаштування real-time парсингу для 3–5 бірж з моніторингом 20–50 пар, reconnect логікою та публікацією в Redis/Kafka займає 1–2 дні. Зв'яжіться з нами для розрахунку вартості. Замовте налаштування вже сьогодні — отримайте консультацію.







