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

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

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

Часто задаваемые вопросы

Последние работы

  • 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 (proxy на опытность) 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. Versioning датасетов критичен — модель обученная на данных за конкретный период должна быть воспроизводима. 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 недели. Свяжитесь с нами для консультации.