Розробка системи реплікації ринкових даних на Kafka

Розробка системи реплікації ринкових даних Ваш торговий бот працює на серверах поруч із біржею. Ризик-менеджмент потребує ті самі дані в дата-центрі за океаном. Проста передача файлів не працює: затримки зростають до секунд, а втрати даних — до 8% тіків при нестабільному каналі. Ми стикалися з ке

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

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

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

  • 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

Розробка системи реплікації ринкових даних

Ваш торговий бот працює на серверах поруч із біржею. Ризик-менеджмент потребує ті самі дані в дата-центрі за океаном. Проста передача файлів не працює: затримки зростають до секунд, а втрати даних — до 8% тіків при нестабільному каналі. Ми стикалися з кейсом, коли клієнт втрачав до 5% тіків через нестабільний канал між Нью-Йорком і Токіо. Потрібна розподілена реплікація в реальному часі, яка витримає сотні тисяч повідомлень на секунду. І не втратить жодного тіка. Реплікація на базі Apache Kafka дозволяє скоротити втрати до нуля та забезпечити консистентність на всіх вузлах. Замовте реплікацію, яка не підведе.

Наша команда має 5+ років досвіду в реплікації ринкових даних, виконала понад 30 успішних проєктів для криптофондів, маркет-мейкерів і проп-трейдингових фірм. Результат: зниження витрат на інфраструктуру до 40% та скорочення часу простою на 80%. Клієнти економлять у середньому $15 000 на місяць на інфраструктурі за рахунок консолідації потоків. Типовий проєкт коштує від $10,000 до $50,000 залежно від складності.

У цій статті розберемо, як побудувати систему реплікації, які топології обрати, як гарантувати консистентність та уникнути типових помилок.

Навіщо потрібна реплікація

Торгова система складається з кількох компонентів, що працюють у різних середовищах:

  • Production trading — co-location поруч із біржею, мінімальна затримка.
  • Research/backtesting — дата-центр з великими обсягами сховища.
  • Risk management — ізольована мережа з обмеженим доступом.
  • Analytics dashboards — доступні широкій аудиторії.

Кожне середовище має отримувати ідентичний потік даних без перевантаження джерела. Реплікація вирішує це завдання, забезпечуючи консистентність, відмовостійкість та масштабованість.

Топології реплікації

Топологія Опис Надійність Затримка
Hub-and-Spoke Один primary-агрегатор збирає дані з бірж, replica-вузли підписуються Низька (single point of failure) Низька
Chain Replication Дані передаються по ланцюжку: біржа → primary → secondary → tertiary Висока (немає SPOF) Висока (кумулятивна)
Pub-Sub (Kafka) Primary пише в Kafka, consumer groups читають незалежно Дуже висока (реплікація топіків) Середня (залежить від mirror)

Для продакшену ми рекомендуємо Pub-Sub через Kafka — це гнучкий варіант, який легко масштабувати. Apache Kafka надає довговічний, масштабований та відмовостійкий обмін повідомленнями. Порівняно з RabbitMQ, Kafka забезпечує в 3 рази вищу пропускну здатність для потоків market data.

Як забезпечити консистентність при реплікації?

Консистентність — ключова проблема при розподіленій реплікації. Ми вирішуємо її комбінацією ідемпотентності consumer'ів та атомарних транзакцій. Consumer'и дедуплікують повідомлення за ключами trade_id або update_id. Для критичних потоків (ризик-менеджмент) використовуємо exactly-once доставку через Kafka Transactions.

from confluent_kafka import Producer producer = Producer({ 'bootstrap.servers': 'kafka:9092', 'enable.idempotence': True, 'transactional.id': 'market-data-producer-1', 'acks': 'all' }) producer.init_transactions() def publish_trade_batch(trades: list[Trade]): producer.begin_transaction() try: for trade in trades: producer.produce( topic=f'market.trades.{trade.exchange}.{trade.symbol}', key=trade.symbol.encode(), value=serialize(trade) ) producer.commit_transaction() except Exception as e: producer.abort_transaction() raise 

Чому Kafka — стандарт для реплікації ринкових даних?

Apache Kafka надає всі необхідні властивості: довговічність (дані зберігаються на диску), масштабованість (горизонтальне партиціонування) та незалежність consumer-груп. Ми налаштовуємо топіки з назвами виду {data_type}.{exchange}.{symbol}.{interval}, що дозволяє легко фільтрувати дані.

Topic: market.trades.binance.BTCUSDT Partition 0: trades (all, ordered by time) Topic: market.orderbook.binance.BTCUSDT Partition 0: snapshots + diffs (ordered by update_id) Topic: market.candles.binance.BTCUSDT.1m Partition 0: 1-minute OHLCV (ordered by candle time) 

Гарантії доставки та крос-датацентр реплікація

У системах market data найчастіше використовується at-least-once delivery: краще отримати дублікат, ніж втратити дані. Споживачі ідемпотентні — дедуплікація за trade_id або update_id. Для risk management та позиційного обліку вмикаємо exactly-once через Kafka Transactions.

Kafka MirrorMaker 2 реплікує топіки між кластерами. Приклад конфігурації MirrorMaker 2:

# mirrormaker2.properties clusters = us-east, eu-west us-east.bootstrap.servers = kafka-us:9092 eu-west.bootstrap.servers = kafka-eu:9092 us-east->eu-west.enabled = true us-east->eu-west.topics = market\.* us-east->eu-west.replication.factor = 2 

EU-кластер отримує репліку всіх market.* топіків із затримкою 50–200 мс для трансатлантичної реплікації. Цього достатньо для більшості аналітичних та ризик-систем.

Управління retention та моніторинг

Market data накопичується швидко. Політики retention:

  • Для tick-даних: 7 днів, після — видаляти.
  • Для daily OHLCV: безстроково, з розмірним лімітом 10 ГБ на партицію.
  • Для order book: log compaction — зберігати лише останній стан на рівень ціни.

Компресія zstd стискає дані на 40–70% без помітного навантаження на CPU. Ключові метрики моніторингу:

Метрика Що показує
Consumer lag Відставання споживачів від producer
Replication latency Затримка між primary та replica кластерами
Producer send rate Швидкість публікації (повідомлень/сек)
Bytes in/out rate Пропускна здатність
Under-replicated partitions Партиції з недостатньою реплікацією

Consumer lag > 5 хвилин для торгового бота — критичний алерт. Для аналітичної системи — warning.

Schema Registry та сумісність форматів

Щоб не зламати споживачів при еволюції схеми, використовуємо Confluent Schema Registry та Avro. Нові опціональні поля з default null — backward compatible зміна.

Приклад схеми Avro
{ "type": "record", "name": "Trade", "namespace": "com.exchange.market", "fields": [ {"name": "exchange", "type": "string"}, {"name": "symbol", "type": "string"}, {"name": "timestamp", "type": "long"}, {"name": "price", "type": {"type": "bytes", "logicalType": "decimal", "precision": 24, "scale": 8}}, {"name": "quantity", "type": {"type": "bytes", "logicalType": "decimal", "precision": 24, "scale": 8}}, {"name": "side", "type": {"type": "enum", "name": "Side", "symbols": ["BUY", "SELL"]}}, {"name": "is_maker", "type": ["null", "boolean"], "default": null} ] } 

Що входить у роботу

  • Архітектурна документація: опис топології, схема потоків, специфікація топіків.
  • Доступи до інфраструктури: налаштування кластерів Kafka, MirrorMaker, Schema Registry.
  • Навчання команди: воркшоп з експлуатації та моніторингу.
  • Підтримка після запуску: допомога в перші 2 тижні продуктивної експлуатації.

Як ми розгортаємо реплікацію: покроковий план

  1. Аналіз вимог: обсяг даних (до 100 000 повідомлень/с), затримки, гарантії доставки, кількість дата-центрів.
  2. Проєктування топології: обираємо між Hub-and-Spoke, Chain або Pub-Sub.
  3. Розгортання Kafka-кластера (від 3 до 7 брокерів) з MirrorMaker 2 та Schema Registry.
  4. Налаштування топіків та retention політик.
  5. Інтеграція consumer'ів з ідемпотентністю та дедуплікацією.
  6. Налаштування моніторингу (consumer lag, latency) та алертів.
  7. Тестування на синтетичних даних та бойових навантаженнях.
  8. Документування та передача команді замовника.

У процесі ми перевіряємо: ідемпотентність producer'ів, налаштування транзакцій, узгодженість retention між кластерами, роботу дедуплікації та відсутність дублів.

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

Базова реалізація займає від 7 до 14 днів: аналітика, розгортання кластера, налаштування MirrorMaker 2 та Schema Registry, моніторинг. Повноцінне рішення з інтеграцією у вашу інфраструктуру — до 4 тижнів. Вартість розраховується індивідуально: залежить від обсягу даних, кількості дата-центрів та необхідних гарантій доставки. Зв'яжіться з нами, щоб оцінити ваш проєкт та отримати консультацію. Замовте розробку системи реплікації під ваші завдання — і ми забезпечимо надійну доставку даних.