Разработка 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-агрегатора для вашей торговой системы.