Розробка pipeline обробки on-chain даних для ML

Ми часто стикаємося із задачею: «хочемо передбачати whale-активність» або «потрібна модель оцінки on-chain кредитного ризику». За цим стоїть інженерна проблема, яку більшість команд недооцінює: сирі блокчейн-дані не придатні для ML-моделей напряму. Структура блоку, raw hex-encoded calldata, адреси в

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

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

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

  • 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

Ми часто стикаємося із задачею: «хочемо передбачати whale-активність» або «потрібна модель оцінки on-chain кредитного ризику». За цим стоїть інженерна проблема, яку більшість команд недооцінює: сирі блокчейн-дані не придатні для ML-моделей напряму. Структура блоку, raw hex-encoded calldata, адреси в bytes20 — це не фічі, це сировина. Між RPC-нодою та навчальною вибіркою лежить кілька тижнів інфраструктурної роботи. Наприклад, щоб побудувати модель відтоку ліквідності з DeFi-протоколу, потрібно зібрати не лише події Transfer, але й внутрішні виклики, trace-інформацію, та нормалізувати часові мітки до єдиного часового поясу. Кожна з цих операцій вимагає окремого пайплайну з контролем помилок та повторюваністю. На практиці без правильного pipeline ви ризикуєте отримати сміттєві фічі, які лише погіршать якість моделі.

Чому сирі блокчейн-дані не придатні для ML?

Raw Transfer лог — це три bytes32 + data bytes. До ML-ознак потрібно пройти:

  • Декодування — ABI-декодинг topics та data
  • Нормалізація адрес — uint256 → checksummed hex, label mapping (біржі, протоколи, MEV-боти)
  • Грошова нормалізація — value / 10^decimals, конвертація в USD через історичний price feed
  • Entity resolution — один EOA може мати сотні транзакцій, але бути одним економічним агентом; смарт-контракти — проксі, реалізації, multisig

Пропуск будь-якого з цих кроків призводить до сміттєвих ознак.

Згідно з документацією Ethereum Foundation, інтеграція з archive node через trace API дозволяє отримати повну історію внутрішніх транзакцій.

Джерела даних: від RPC до спеціалізованих провайдерів

Публічні RPC (eth_getLogs, eth_getBlockByNumber) — найдоступніше, але найменш придатне для ML джерело. Їх обмеження: rate limits (Infura/Alchemy — 10-333 req/s на платних тарифах), відсутність internal transactions без trace_ namespace, відсутність pre/post state без архівної ноди. Archive node з trace API дає повну історію, але вимагає Erigon з дисковим простором ~2.5 TB для Ethereum mainnet та синхронізацією 3-5 днів. Формати trace_ різняться між Erigon та Geth/Besu — парсер доводиться адаптувати. Firehose (StreamingFast/The Graph) експортує кожен блок з деревом викликів та state diffs за <500ms, забезпечуючи швидкість 100k+ блоків за хвилину — в 20-100 разів швидше за RPC. Спеціалізовані постачальники (Nansen, Dune, Flipside, Allium) дають готові нормалізовані таблиці, але з затримкою оновлення 1-24 години та обмеженим контролем над схемою. Для production ML оптимально комбінувати: Firehose для історичного завантаження та archive node для real-time стрімінгу.

Як гарантується point-in-time коректність на рівні backend?

Це ключова проблема. Ознаки повинні бути обчислені лише з даних, доступних до моменту передбачення. Типова помилка: використання total_tx_count адреси замість tx_count_at_time_T. Патерн: point-in-time correct features. Кожен рядок у feature store має entity_id, feature_timestamp, feature_value. При генерації навчальної вибірки джойн йде по entity_id та feature_timestamp <= label_timestamp.

-- Point-in-time join SELECT l.wallet_address, l.label, l.label_timestamp, f.tx_count, f.unique_contracts, f.volume_usd_30d FROM labels l ASOF JOIN wallet_features f ON l.wallet_address = f.wallet_address AND f.feature_timestamp <= l.label_timestamp 

ASOF JOIN — нативна операція в ClickHouse та TimescaleDB, в PostgreSQL емулюється через LATERAL.

Offline store — історичні фічі для навчання. ClickHouse або Parquet на S3 з Hive-partitioning за датою. Online store — актуальні фічі для inference. Redis Hash structures: HGETALL wallet:{address}:features. Оновлюється при кожному новому блоці для активних адрес.

Архітектура production pipeline

Шар інжестії

Рекомендована архітектура — event-driven з розділенням hot та cold path:

[Archive Node / Firehose] ↓ [Kafka / Redpanda] ← hot path: < 1s latency ↓ [Stream Processor] ← Flink або кастомний consumer / \ [Raw Store] [Feature Store] ← cold: S3/Parquet, hot: Redis/Feast 

Kafka topic per chain, ключ = block_number:log_index. Це гарантує порядок і дозволяє replay при помилках обробки. Retention залежить від задачі: для real-time фіч — 7 днів, для перенавчання — повний архів в S3. Для Ethereum mainnet: ~6000 транзакцій/блок × ~6500 блоків/день = ~39M транзакцій/день. При середньому розмірі транзакції з trace ~2KB — ~75GB/день сирих даних. Плануйте сховище.

Feature engineering та сховища

Це найтрудомісткіша частина. Типові on-chain ознаки для різних ML-задач: Wallet profiling (DeFi credit scoring, Sybil detection):

Ознака Джерело Складність
Вік адреси (блоки з першої TX) eth_getTransactionCount history низька
Унікальні контракти взаємодії event logs середня
Gas percentile (проксі на досвідченість) TX history низька
Час між транзакціями (ритмічність) TX timestamps середня
Nonce gaps (втрачені TX) nonce vs tx count середня
DeFi protocol diversity contract label mapping висока
Liquidation history protocol-specific events висока

MEV detection:

  • Sandwich attack pattern: три TX в одному блоці, одна адреса, оточують target TX
  • Arbitrage: циклічні трансфери токенів, що повертаються до sender в рамках однієї TX
  • Flashloan: FlashLoan event + position delta = 0 до кінця блоку

Whale activity prediction:

  • Великі трансфери з exchange deposit addresses → ймовірність sell pressure
  • Accumulation pattern: множинні невеликі покупки з різних адрес → один отримувач
# Приклад feature engineering для wallet scoring import polars as pl def compute_wallet_features(txs: pl.DataFrame) -> pl.DataFrame: return txs.group_by("from_address").agg([ pl.col("block_number").min().alias("first_seen_block"), pl.col("block_number").max().alias("last_seen_block"), pl.count("hash").alias("tx_count"), pl.col("to_address").n_unique().alias("unique_contracts"), pl.col("gas_price").quantile(0.5).alias("gas_price_median"), pl.col("value_usd").sum().alias("total_volume_usd"), pl.col("block_timestamp").diff().dt.total_seconds() .mean().alias("avg_interval_seconds"), ]) 

Polars замість Pandas — різниця у швидкості обробки великих датасетів (мільйони рядків) становить 5-20x.

Обробка реорганізацій та MLOps

Реорги на рівні ML-даних — серйозна проблема. Якщо фічі обчислені з блоку, який згодом став orphaned, навчальна вибірка містить нереальні дані. Рішення:

  • Confirmation lag — індексувати лише блоки старші за N блоків (зазвичай 12-32 для фінальності на PoS Ethereum). Додає затримку, але усуває проблему.
  • Versioned features — зберігати (entity, block_hash, features), при реорзі помічати orphaned записи. Складніше, але дозволяє працювати з малою затримкою.

MLOps інтеграція. Pipeline повинен дружити з існуючим ML-стеком: Feature generation → навчання: експорт в Parquet/CSV для DVC або MLflow artifacts. Версіонування датасетів критичне — модель навчена на даних за конкретний період повинна бути відтворювана. Inference pipeline: новий блок → обчислення дельта-фіч → update в online store → тригер inference. Latency бюджет зазвичай 1-10 секунд від блоку до передбачення. Model drift monitoring: on-chain дані змінюються структурно (merge, нові протоколи, зміни патернів використання). Потрібен моніторинг дистрибуції вхідних ознак — Evidently AI або кастомний.

Порівняння джерел даних
Характеристика Firehose Public RPC Archive Node (Erigon)
Швидкість 100k+ блоків/хв 1-5k блоків/хв 5-20 блоків/хв
Latency <500ms 1-3s 2-5s
Витрати (self-hosted) Високі Низькі Середні
Повнота даних Повний trace Тільки зовнішні TX Повний trace + state

Типові етапи проекту

Data audit (1-2 тижні)

Визначення потрібних сигналів, їх джерел, доступності історичних даних. Прототип інжестора на невеликому блок-діапазоні.

Historical backfill (2-4 тижні)

Завантаження історичних даних, нормалізація, label mapping. Найтрудомісткіший етап.

Feature pipeline (2-3 тижні)

Реалізація feature engineering, point-in-time logic, сховища.

Real-time path (1-2 тижні)

Стрімінг з ноди, online store, inference інтеграція.

MLOps (1-2 тижні)

Моніторинг дрейфу, версіонування датасетів, автоматизація перенавчання.

Разом: 7-13 тижнів до production-ready pipeline. Оцінка сильно залежить від кількості ланцюгів, глибини історичних даних та вимог до latency inference.

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

  • Підготовка архітектури pipeline під вашу задачу
  • Налаштування інфраструктури (Kafka, ClickHouse, Redis)
  • Розробка feature engineering для цільових ознак
  • Інтеграція з MLOps (MLflow, DVC)
  • Документація та навчання команди

Наша команда має багаторічний досвід у блокчейн-розробці та реалізувала більше 20 проектів з аналізу on-chain даних. Використання нашого pipeline дозволяє скоротити витрати на інфраструктуру до 40% порівняно з самостійною розробкою. Замовте аудит ваших on-chain даних — отримайте прототип pipeline за 2 тижні. Зв'яжіться з нами для консультації.