Мы разрабатываем data lake для блокчейн-данных — слой, который решает главную проблему on-chain аналитики: сырые данные блокчейна не приспособлены для сложных запросов. JSON-RPC узлы отвечают на вопрос «что произошло в блоке X», но не на «покажи все свопы Uniswap V3 за последние 30 дней по адресам с объёмом > 10k USDC». Data lake превращает сырые блоки в структурированные, индексированные и быстро запрашиваемые таблицы.
Ethereum mainnet сегодня — это около 20 миллионов блоков, примерно 2 миллиарда транзакций и терабайты event logs. Полная история Ethereum в формате Parquet занимает 3–4 TB. На каждый новый блок (каждые 12 секунд) добавляются сотни транзакций и тысячи log-записей. К этому добавляются BSC, Polygon, Arbitrum, Base — у каждой сети своя история и скорость роста. Наш data lake объединяет их в единую аналитическую среду.
Почему data lake необходим для аналитики on-chain?
Три класса данных с разными характеристиками:
Blocks & transactions — структурированные, предсказуемая схема. Главная сложность: reorgs — временные форки, после которых цепочка переписывается. Пайплайн должен уметь откатывать уже записанные данные.
Event logs — самые ценные для аналитики. Transfer, Swap, Liquidation, Mint — всё это EVM events. Проблема: декодирование ABI. Без ABI контракта log — это просто байты с topics. Мы создаём реестр ABI через Etherscan API и Sourcify, чтобы декодировать миллионы событий автоматически.
Traces (internal transactions) — вызовы между контрактами, не создающие прямой транзакции. Без трейсов невидима значительная часть DeFi: flash loan внутри одной tx, recursive liquidations, MEV bundle. Получение трейсов через debug_traceTransaction — тяжёлая операция, доступна только на archive-узлах.
Как обрабатывать reorgs?
Reorg — главная головная боль любого блокчейн data pipeline. Ethereum с Proof-of-Stake имеет probabilistic finality через несколько блоков и полный finality через ~12.8 минут (2 эпохи). L2 сети имеют ещё более сложную модель.
Стандартный подход:
- Записывать блоки с confirmmation lag (ждать N подтверждений перед записью в финальный слой). Для Ethereum: 32–64 блока.
- Хранить staging-слой для последних M блоков — данные туда пишутся немедленно, но помечаются как
pending. - Подписка на
Reorganizationсобытия от ноды (WebSocketnewHeads+ сравнение parentHash). При reorg — удаляем затронутые блоки из staging и переприменяем новую цепочку.
Для Iceberg это элегантно решается через time travel и merge операции. Для ClickHouse — через ReplacingMergeTree с version-столбцом.
Архитектура data lake
Слой ingestion
Два подхода:
- Node-based ingestion — прямое подключение к узлу через WebSocket. Подписка на новые блоки + backfill через
eth_getLogsbatch calls. Требует archive node. Для backfill миллионов блоков используем параллельную обработку с asyncio. - Third-party data providers — Goldsky, Envio, Substreams. Быстрее старта, но vendor lock-in и дороже на масштабе.
Хранилище: выбор формата и движка
Для сырых данных блокчейна оптимален columnar storage:
- Apache Parquet на S3/GCS — стандарт. Компрессия zstd уменьшает объём в 5–10x. Партиционирование по дате и номеру блока.
- Apache Iceberg поверх Parquet — ACID, schema evolution, time travel. Критично для reorgs.
- ClickHouse — OLAP для горячих запросов. Сотни миллионов строк за секунды.
Типичная двуслойная архитектура:
Raw layer (S3 + Parquet/Iceberg) ↓ ETL (dbt / Spark / Flink) Serving layer (ClickHouse / BigQuery) ↓ Query API Analytics / Trading systems / Dashboards Декодирование ABI и enrichment
Сырые event logs содержат topics (хэши event signatures) и data (ABI-encoded). Для декодирования нужен ABI реестр:
from eth_abi import decode from web3 import Web3 TRANSFER_TOPIC = Web3.keccak(text="Transfer(address,address,uint256)").hex() def decode_transfer(log: dict) -> dict | None: if log["topics"][0] != TRANSFER_TOPIC: return None from_addr = "0x" + log["topics"][1][-40:] to_addr = "0x" + log["topics"][2][-40:] amount = decode(["uint256"], bytes.fromhex(log["data"][2:]))[0] return {"from": from_addr, "to": to_addr, "amount": amount} Для массового декодирования создаём реестр ABI — таблицу с маппингом contract_address → ABI. Источники: Etherscan API, Sourcify, 4byte.directory. Неизвестные контракты обрабатываем как raw bytes, enrichment — по мере появления ABI.
Token metadata enrichment: для ERC-20 трансферов нужны decimals, symbol, цена. Цены берём из Uniswap V3 TWAP записей или внешних API (исторические данные).
Схема данных и ключевые таблицы
CREATE TABLE decoded_events ( block_number UInt64, block_timestamp DateTime, tx_hash FixedString(66), log_index UInt32, contract FixedString(42), event_name LowCardinality(String), chain_id UInt32, params String, -- JSON INDEX idx_contract (contract) TYPE bloom_filter GRANULARITY 4, INDEX idx_event (event_name) TYPE set(100) GRANULARITY 4 ) ENGINE = ReplacingMergeTree(block_number) PARTITION BY toYYYYMM(block_timestamp) ORDER BY (chain_id, contract, block_number, log_index); Отдельные таблицы для высокочастотных event types: erc20_transfers, uniswap_v3_swaps, aave_liquidations. Партиционирование по месяцам.
Пример партиционирования и оптимизации
Для сети Ethereum партиционируем таблицу событий по месяцам. Это позволяет быстро удалять устаревшие данные и эффективно сканировать временные диапазоны. Индексы bloom_filter на контракте ускоряют фильтрацию по адресам.Пример запроса:
SELECT sum(amount) FROM erc20_transfers WHERE contract = '0xdAC17F958D2ee523a2206206994597C13D831ec7' AND block_timestamp >= '2023-01-01' выполняется за доли секунды на 50M строк.
Сравнение способов получения данных
| Параметр | Node-based | Third-party (Goldsky) |
|---|---|---|
| Скорость старта | Средняя (настройка ноды) | Высокая (API ключ) |
| Контроль данных | Полный | Ограниченный вендором |
| Стоимость на масштабе | Низкая (свои ноды) | Высокая (плата за объём) |
| Обработка reorgs | Собственный механизм | Встроенная (но непрозрачно) |
Что входит в результат
На выходе вы получаете:
- Рабочий data lake с выбранными сетями и событиями.
- Документацию схемы данных и ETL-пайплайна.
- Доступ к ClickHouse (или другому serving layer) с примерами запросов.
- Мониторинг lag’а и алерты при проблемах.
- Обучение команды и 1 месяц поддержки.
Возможность расширения: добавление новых контрактов или сетей — подключается через конфигурацию без изменения кода.
Этапы разработки
| Фаза | Содержание | Длительность |
|---|---|---|
| Design | Определение scope, сетей/событий, схема данных | 1–2 нед |
| Core ingestion | WebSocket listener, backfill, reorg handler | 3–4 нед |
| ABI registry | Накопление ABI, декодирование, enrichment | 2–3 нед |
| Storage layer | Parquet/Iceberg, ClickHouse, ETL | 3–4 нед |
| Serving API | REST/GraphQL, rate limiting | 2–3 нед |
| Monitoring & ops | Airflow, алерты, документация | 1–2 нед |
Почему стоит работать с нами
Наш опыт — более 5 лет в блокчейн-инженерии, реализовано 20+ data pipeline для DeFi-протоколов и крипто-фондов. Мы разрабатываем под ключ — от проектирования схемы до деплоя и мониторинга. Оценим ваш проект за 2 дня — свяжитесь с нами для консультации. Закажите разработку data lake под ключ, чтобы ускорить аналитику on-chain-данных.
Trust-слова: гарантия качества, сертифицированные специалисты, большой опыт работы с L1/L2. Ссылка: Подробнее о структуре блокчейна можно прочитать на Wikipedia: Блокчейн.







