Мы часто сталкиваемся с задачей: «хотим предсказывать 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 недели. Свяжитесь с нами для консультации.







