Торгові боти, маркет-мейкери та ризик-системи залежать від свіжих ринкових даних — затримка в мілісекунди може коштувати прибутку. В одному проєкті клієнт втратив десятки тисяч доларів за місяць через 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-агрегатора для вашої торгової системи.







