Розробка системи реплікації ринкових даних
Ваш торговий бот працює на серверах поруч із біржею. Ризик-менеджмент потребує ті самі дані в дата-центрі за океаном. Проста передача файлів не працює: затримки зростають до секунд, а втрати даних — до 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 тижні продуктивної експлуатації.
Як ми розгортаємо реплікацію: покроковий план
- Аналіз вимог: обсяг даних (до 100 000 повідомлень/с), затримки, гарантії доставки, кількість дата-центрів.
- Проєктування топології: обираємо між Hub-and-Spoke, Chain або Pub-Sub.
- Розгортання Kafka-кластера (від 3 до 7 брокерів) з MirrorMaker 2 та Schema Registry.
- Налаштування топіків та retention політик.
- Інтеграція consumer'ів з ідемпотентністю та дедуплікацією.
- Налаштування моніторингу (consumer lag, latency) та алертів.
- Тестування на синтетичних даних та бойових навантаженнях.
- Документування та передача команді замовника.
У процесі ми перевіряємо: ідемпотентність producer'ів, налаштування транзакцій, узгодженість retention між кластерами, роботу дедуплікації та відсутність дублів.
Терміни та вартість
Базова реалізація займає від 7 до 14 днів: аналітика, розгортання кластера, налаштування MirrorMaker 2 та Schema Registry, моніторинг. Повноцінне рішення з інтеграцією у вашу інфраструктуру — до 4 тижнів. Вартість розраховується індивідуально: залежить від обсягу даних, кількості дата-центрів та необхідних гарантій доставки. Зв'яжіться з нами, щоб оцінити ваш проєкт та отримати консультацію. Замовте розробку системи реплікації під ваші завдання — і ми забезпечимо надійну доставку даних.







