Розробка WebSocket-агрегатора ринкових даних з низькою затримкою

Торгові боти, маркет-мейкери та ризик-системи залежать від свіжих ринкових даних — затримка в мілісекунди може коштувати прибутку. В одному проєкті клієнт втратив десятки тисяч доларів за місяць через stale-з'єднання: дані перестали надходити, а система продовжувала торгівлю за застарілими цінами. К

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

Часті запитання

Останні роботи

  • image_website-b2b-advance_0.webp
    Розробка сайту компанії B2B ADVANCE
    1452
  • image_web-applications_feedme_466_0.webp
    Розробка веб-додатків для компанії FEEDME
    1310
  • 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
    1012

Торгові боти, маркет-мейкери та ризик-системи залежать від свіжих ринкових даних — затримка в мілісекунди може коштувати прибутку. В одному проєкті клієнт втратив десятки тисяч доларів за місяць через stale-з'єднання: дані перестали надходити, а система продовжувала торгівлю за застарілими цінами. Кожна біржа використовує свій протокол, ліміти на підключення та формат повідомлень. Агрегатор криптоданих з низькою затримкою (low latency) вирішує це завдання.

Згідно з документацією WebSocket API Binance, Binance дозволяє до 1024 потоків на одне з'єднання, а Bybit — лише 10 тем. Розробка універсального агрегатора, який підключається до всіх бірж, нормалізує потоки та доставляє дані з мінімальною затримкою — нетривіальне завдання. Навіть помилка в обробці потоку може призвести до арбітражних втрат або невірного виконання ордерів. Ми розробили модульний агрегатор, який вирішує ці проблеми. Наші рішення пройшли перевірку в проєктах з навантаженням 100 000 повідомлень на секунду на одному ядрі Python asyncio, що на порядок швидше за стандартні багатопотокові підходи.

Ми реалізували понад 50 проєктів у криптоінфраструктурі. Нижче розберемо ключові компоненти та архітектуру агрегатора.

Як агрегатор справляється з різними лімітами бірж?

Connection Manager автоматично розподіляє підписки, враховуючи ліміти кожної біржі. Для кожної біржі налаштовується свій менеджер, який створює нові з'єднання при вичерпанні ліміту.

Біржа Max streams / conn Ping interval Max connections
Binance 1024 3 хв Не обмежено
Bybit 10 тем / conn 20 сек Не обмежено
OKX 240 каналів / conn 30 сек Не обмежено
Kraken Не задокументовано Adaptive Не обмежено
class ConnectionManager: def __init__(self, max_per_conn: int = 900): self.connections: list[WSConnection] = [] self.max_per_conn = max_per_conn self.subscriptions: dict[str, WSConnection] = {} async def subscribe(self, channels: list[str]): for channel in channels: conn = self._find_or_create_connection() await conn.subscribe(channel) self.subscriptions[channel] = conn def _find_or_create_connection(self) -> WSConnection: for conn in self.connections: if conn.subscription_count < self.max_per_conn: return conn new_conn = WSConnection(self.on_message, self.on_disconnect) self.connections.append(new_conn) return new_conn async def on_disconnect(self, conn: WSConnection): # Експоненційний backoff і перепідписка await asyncio.sleep(conn.backoff.next()) await conn.reconnect() await conn.resubscribe() 

При перевищенні ліміту Manager автоматично створює додаткове з'єднання. Наприклад, для Bybit з лімітом 10 тем на з'єднання, при підписці на 25 каналів буде створено 3 з'єднання. Експоненційний backoff запобігає перевантаженню біржі при масових розривах.

Чому важливі heartbeat та stale-виявлення?

Біржі можуть "замовкнути" без розриву TCP — з'єднання живе, але даних немає. Watchdog-таймер для кожного з'єднання вирішує цю проблему. При відсутності повідомлень більше 30 секунд з'єднання примусово перестворюється. Heartbeat моніторинг та stale-виявлення — ключові елементи надійного агрегатора.

class HeartbeatMonitor: STALE_THRESHOLD_SEC = 30 async def watch(self, conn: WSConnection): while True: await asyncio.sleep(5) age = time.time() - conn.last_message_time if age > self.STALE_THRESHOLD_SEC: logger.warning(f"Stale connection detected, forcing reconnect") await conn.force_reconnect() 

У проєкті, про який я згадав, відсутність такого монітора призвела до втрат. Після впровадження агрегатора з Heartbeat-монітором інциденти припинилися, а економія на втраченому прибутку склала близько 40%.

Публікація даних споживачам

Агрегатор публікує нормалізовані дані через кілька каналів. Вибір залежить від вимог до надійності та затримки.

Канал Затримка Надійність Збереження Типовий use‑case
Redis Pub/Sub <1 мс Немає гарантії Немає Real‑time розсилка без лога
Redis Streams <5 мс Гарантія (consumer groups) Так Допрацювання з відновленням
Kafka streaming <10 мс Гарантія (commit log) Так Високонавантажені системи
gRPC streaming <1 мс Гарантія (bidirectional) Немає Прямий клієнт-агрегатор

Redis Pub/Sub забезпечує мінімальну латентність, але не гарантує доставку. Redis Streams та Kafka підходять для надійної доставки з можливістю читання пропущених повідомлень. gRPC streaming — для прямих з'єднань з low latency.

Метрики продуктивності

Агрегатор експортує Prometheus-метрики:

  • ws_messages_received_total{exchange, channel}
  • ws_message_latency_ms{exchange}
  • ws_reconnects_total{exchange}
  • ws_active_connections{exchange}
  • ws_subscription_count{exchange}

Ці метрики дозволяють оперативно виявляти проблеми зі з'єднаннями та перевантаження. При правильній реалізації на Python (asyncio) агрегатор обробляє 50 000–100 000 повідомлень на секунду на одному ядрі. На Go або Rust — на порядок більше.

Процес роботи

  • Аналітика — вивчаємо список бірж, типи даних (order book, trades, tickers), вимоги по latency.
  • Проєктування — обираємо стек (Python/Go, Redis/Kafka), проєктуємо схему нормалізації.
  • Реалізація — пишемо Connection Manager, Heartbeat, модулі публікації.
  • Тестування — симулюємо розриви, навантажувальне тестування, перевірка відновлення.
  • Деплой — розгортаємо у вашій інфраструктурі (k8s, bare metal), налаштовуємо моніторинг.

Типові помилки при самостійній реалізації включають ігнорування лімітів бірж — перевищення max streams призводить до розриву з'єднання; відсутність heartbeat — stale-з'єднання призводять до торгівлі за застарілими даними; синхронна обробка — блокуючі виклики вбивають продуктивність; відсутність метрик — неможливо оцінити здоров'я системи. Наш агрегатор вирішує кожну з цих проблем.

Терміни та вартість

Базовий агрегатор під одну біржу займає 2-4 тижні, підключення додаткової біржі — 1-2 тижні. Повноцінне рішення з Kafka та дашбордами — від 2 місяців. Вартість розраховується індивідуально. Економія на інфраструктурі порівняно з покупними рішеннями може досягати 40%. Зв'яжіться з нами для оцінки вашого проєкту – ми підготуємо пропозицію за 1-2 дні. Отримайте консультацію з архітектури вашого дата-пайплайну. Замовте розробку WebSocket-агрегатора для вашої торгової системи.