Розробка потокової обробки блокчейн-даних (Kafka/Flink)

Розробка потокової обробки блокчейн-даних (Kafka/Flink) Ми часто стикаємося з ситуацією, коли нода Ethereum в режимі real-time генерує близько 2–5 МБ даних на секунду в періоди високої активності мережі. Це події Transfer, виклики контрактів, зміни стану. Якщо ваша аналітична система або торговий

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

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

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

  • 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

Розробка потокової обробки блокчейн-даних (Kafka/Flink)

Ми часто стикаємося з ситуацією, коли нода Ethereum в режимі real-time генерує близько 2–5 МБ даних на секунду в періоди високої активності мережі. Це події Transfer, виклики контрактів, зміни стану. Якщо ваша аналітична система або торговий двигун отримують ці дані через періодичний polling RPC-ноди — ви працюєте з застарілими даними та пропускаєте події. Для задач, де затримка в 1–2 блоки критична (арбітраж, liquidation monitoring, fraud detection), потрібна потокова архітектура з гарантіями доставки. Наші інженери будують такі системи під ключ — від топіків до дашбордів.

Чому потокова обробка блокчейн-даних критична для DeFi?

Арбітражні боти, liquidation моніторинг та виявлення MEV потребують затримки менше 500 мс від появи блоку до прийняття рішення. Polling RPC-ноди через JSON-RPC дає затримку в кілька секунд і не гарантує доставку всіх подій. Потокова архітектура на Kafka забезпечує збереження всіх даних з можливістю перечитування, а Flink дозволяє виконувати ковзні агрегації та детектувати складні патерни в реальному часі.

Джерела даних: від ноди до Kafka

WebSocket підписки vs polling

Стандартний eth_subscribe("newHeads") через WebSocket дає сповіщення про новий блок без затримки polling. Але WebSocket-з'єднання нестабільне на довгих періодах — потрібен reconnect з catchup логікою:

func (s *NodeSubscriber) subscribeWithRecovery(ctx context.Context) error { for { lastBlock, _ := s.db.GetLastProcessedBlock() // Догнати пропущені блоки при reconnect if err := s.catchUpFromBlock(ctx, lastBlock+1); err != nil { return err } // Підписатися на нові блоки sub, err := s.client.SubscribeNewHead(ctx, s.headers) if err != nil { time.Sleep(backoffDuration) continue } select { case err := <-sub.Err(): log.Warnf("subscription error: %v, reconnecting", err) case <-ctx.Done(): return nil } } } 

Firehose protocol (StreamingFast/Pinax)

Для Ethereum та інших EVM-мереж найефективніший спосіб отримання raw даних — Firehose (StreamingFast), що інструментує ноду на рівні бінарника та експортує блоки в protobuf з мінімальною затримкою. Throughput на порядок вище ніж через JSON-RPC. Для проектів з вимогою повної історичної прокрутки — Firehose + зберігання flat files в S3/GCS дозволяє відтворювати будь-який діапазон блоків без повторної синхронізації ноди.

Kafka як транспортний шар

Kafka — черга зі збереженням (log-based). На відміну від RabbitMQ/Redis Streams, Kafka зберігає всі повідомлення в налаштований retention період (дні, тижні), що дозволяє consumers перечитувати дані. Це критично для блокчейн-аналітики: нова consumer group може прочитати всю історію подій без звернення до ноди.

Топології топіків для blockchain pipeline:

raw.blocks → сирі блоки (partitioned by block_number % N) raw.transactions → всі транзакції raw.logs → всі event logs decoded.transfers → декодовані ERC-20 Transfer події decoded.swaps → декодовані Swap події (Uniswap, Curve, etc.) alerts.large-txns → транзакції > threshold analytics.prices → агреговані цінові дані 

Partitioning strategy важлива: для подій конкретного контракту — partition by contractAddress (гарантує ordering). Для транзакцій — partition by from address або blockNumber.

Apache Flink: stateful stream processing

Flink — правильний інструмент для задач, які потребують стану: ковзні агрегати, join потоків, виявлення патернів у часі. Spark Streaming — батчинг під виглядом стрімінгу (micro-batches). Flink — справжній event-time processing.

Декодування ABI on-the-fly

Вхідні logs — сирі hex дані. Flink job повинен декодувати їх у типізовані події:

public class LogDecoderFunction extends RichFlatMapFunction<RawLog, DecodedEvent> { private Map<String, ContractABI> abiRegistry; @Override public void flatMap(RawLog log, Collector<DecodedEvent> out) { String contractAddress = log.getAddress().toLowerCase(); ContractABI abi = abiRegistry.get(contractAddress); if (abi == null) return; // невідомий контракт String topic0 = log.getTopics().get(0); EventDefinition eventDef = abi.findEventBySignatureHash(topic0); if (eventDef != null) { DecodedEvent decoded = AbiDecoder.decode(eventDef, log); out.collect(decoded); } } } 

ABI реєстр завантажується з PostgreSQL/Redis при старті job та оновлюється через Broadcast State pattern — без перезапуску job при додаванні нових контрактів.

Часові вікна та агрегації

Задача: обчислювати 5-хвилинний VWAP (Volume Weighted Average Price) по свапах Uniswap V3 в режимі реального часу.

DataStream<SwapEvent> swaps = source .filter(e -> e.getType().equals("Swap")) .map(e -> (SwapEvent) e); DataStream<VWAPResult> vwap = swaps .keyBy(SwapEvent::getPoolAddress) .window(TumblingEventTimeWindows.of(Time.minutes(5))) .aggregate(new VWAPAggregator(), new VWAPWindowFunction()); 

Event time vs processing time — принциповий вибір. Event time (час блоку) дає детерміновані результати при переграванні історії. Processing time швидше, але дає різні результати при replay.

Watermarks для обробки late events — блокчейн транзакції можуть приходити в Kafka з невеликою затримкою:

WatermarkStrategy.<RawLog>forBoundedOutOfOrderness(Duration.ofSeconds(10)) .withTimestampAssigner((log, ts) -> log.getBlockTimestamp() * 1000L) 

Складні патерни: CEP для виявлення аномалій

Flink CEP (Complex Event Processing) дозволяє описувати послідовності подій. Задача: детектувати sandwich attack — front-run транзакція, жертва, back-run транзакція в одному блоці.

Pattern<DecodedEvent, ?> sandwichPattern = Pattern .<DecodedEvent>begin("frontrun") .where(e -> e.isSwap() && e.getGasPrice() > threshold) .next("victim") .where(e -> e.isSwap() && samePool(e, "frontrun")) .next("backrun") .where(e -> e.isSwap() && samePool(e, "frontrun") && e.getSender().equals(frontrunSender(e))) .within(Time.seconds(12)); // в межах одного блоку 

State backend та відмовостійкість

Як ми забезпечуємо exactly-once доставку?

Flink checkpoint — знімок стану всіх операторів в S3/HDFS. При збої — відновлення з останнього checkpoint, Kafka consumer offset зберігається атомарно з state. Це гарантує exactly-once семантику для більшості операторів.

RocksDB state backend — обов'язковий для production при великому стані (мільйони ключів). In-memory backend не масштабується.

Деталі про checkpointing Checkpointing інтервал 60 секунд забезпечує баланс між продуктивністю та відновленням. При збої відновлення займає не більше 2 хвилин.

Моніторинг та dead letter queues

Необроблені події (невідомий ABI, помилка парсингу, неочікуваний формат) не можна просто дропати. Dead letter queue (DLQ) в окремий Kafka топік зі збереженням оригінального повідомлення та stack trace помилки — стандартний патерн.

Метрики Flink + Prometheus + Grafana: lag по кожному топіку, throughput операторів, backpressure по графу job. Backpressure — перший індикатор, що downstream не справляється.

Типові use cases та затримки

Use case Допустима затримка Інструмент
MEV bot / арбітраж < 100мс WebSocket → in-process
Liquidation monitoring < 1 сек Kafka + Flink CEP
DeFi аналітика real-time 1–5 сек Kafka + Flink aggregations
Ончейн аналітика/BI < 1 хв Kafka + Flink → ClickHouse
Історичний аналіз без обмежень Firehose → S3 → Spark/dbt

Порівняння інструментів потокової обробки

Інструмент Підхід Гарантія доставки Затримка
Apache Flink True streaming, event-time Exactly-once < 100 мс
Kafka Streams Stream-table duality At-least-once < 100 мс
Spark Streaming Micro-batches Exactly-once (via checkpoint) ~ 1 сек
Akka Streams Reactive streams At-most-once < 50 мс

Інфраструктура та стек

Мінімальний production кластер: 3 Kafka брокери (3 репліки для durability), Flink cluster з 1 JobManager + 3–5 TaskManager pods в Kubernetes. Сховище результатів — ClickHouse для аналітичних запитів (колонкове, швидкі aggregations на великих обсягах) або PostgreSQL + TimescaleDB для метрик часових рядів.

Керовані сервіси скорочують операційне навантаження: Confluent Cloud (Kafka), Amazon Kinesis (альтернатива для AWS-native стека). Для on-premise або compliance вимог — власний кластер.

Що входить в розробку системи

  • Архітектура потокового пайплайну від джерел до сховищ
  • Налаштування Kafka: топіки, партиціонування, retention політики
  • Розробка Flink jobs: декодування ABI, агрегації, CEP патерни
  • Моніторинг та оповіщення: Prometheus + Grafana дашборди
  • Документація та навчання команди
  • Підтримка після запуску (згідно SLA)

Наша команда має 7+ років досвіду в розробці високонавантажених систем для Crypto та DeFi, реалізувала 30+ проектів. Готові оцінити ваш проект — напишіть нам. Оцінка займає 2 робочі дні.