На третій день після запуску аналітики по Uniswap v3 виявляєш, що eth_getLogs із широким фільтром починає таймаутити, агрегати розходяться через пропущені реорганізації, а твій PostgreSQL пухне з гігабайтними таблицями без партиціонування. On-chain ETL — це не просто "читаємо логи і пишемо в базу". Це система з гарантіями консистентності, обробкою реорганізацій, трансформацією даних і керованим беклогом. Ми будуємо її правильно з першого разу, і в цьому тексті розберемо ключові архітектурні рішення.
Приклад із практики: один із наших клієнтів втратив 2 тижні на відновлення даних через неправильну обробку реоргів. Після впровадження нашого пайплайну економія часу склала 40% на історичній синхронізації, а вартість володіння інфраструктурою знизилась на 30% за рахунок оптимізації зберігання.
ETL без обробки реорганізацій — це не ETL, а генерація сміття. Наш досвід говорить: 90% проблем вирішуються правильною архітектурою на старті.
Що таке on-chain ETL-пайплайн і навіщо він потрібен?
On-chain ETL-пайплайн — це система, яка видобуває сирі дані з блокчейну (логи подій, внутрішні транзакції, зміни стану), трансформує їх у структуровані записи (декодування ABI, збагачення цінами, нормалізація сум) і завантажує в аналітичне сховище. Без такого пайплайну неможливо будувати дашборди DeFi-протоколів, відстежувати ліквідність у реальному часі чи проводити історичний аналіз. Основні складнощі: реорганізації ланцюжків, величезні обсяги (до 15M+ блоків в Ethereum) і необхідність гарантувати консистентність при паралельній інгестії.
Як влаштована архітектура: три шари ETL
Класичний ETL (Extract — Transform — Load) у блокчейн-контексті набуває специфіки: джерело даних immutable, але не фінальне (реорги), обсяги вимірюються сотнями мільйонів подій, а latency може бути як секунди, так і години — залежить від задачі.
Extract: інгестія з ноди
Вибір джерела даних визначає все інше. Три рівні з наростаючою складністю:
- Logs/Events — те, що контракт явно емітує. Дешево, швидко, структуровано через ABI. Обмеження: тільки те, що розробник вирішив логувати.
- Traces (internal transactions) — всі виклики всередині транзакції, включаючи ETH-перекази без подій. Вимагає
debug_traceTransactionабоtrace_block(Parity-style). Не всі ноди підтримують; Erigon — найкращий вибір для trace-heavy задач. - State diffs — зміни storage slot-ів по блоку. Максимальна повнота, але величезний обсяг даних і складність інтерпретації без ABI.
Для більшості DeFi-задач достатньо logs + traces. State diffs потрібні для MEV-аналітики та моніторингу контрактів без подій (наприклад, legacy WETH).
Паттерни отримання даних:
# Polling з експоненційним backoff async def fetch_logs_range( rpc: AsyncWeb3, from_block: int, to_block: int, addresses: list[str], topics: list[str], ) -> list[Log]: try: return await rpc.eth.get_logs({ "fromBlock": from_block, "toBlock": to_block, "address": addresses, "topics": [topics], }) except ValueError as e: # "Log response size exceeded" — ділимо діапазон навпіл if "exceeded" in str(e) and from_block < to_block: mid = (from_block + to_block) // 2 left = await fetch_logs_range(rpc, from_block, mid, addresses, topics) right = await fetch_logs_range(rpc, mid + 1, to_block, addresses, topics) return left + right raise Цей паттерн рекурсивного бісекта — обов'язковий. Публічні RPC (і навіть Alchemy/Infura) ріжуть відповіді за розміром. Без нього пайплайн буде падати на активних блоках.
WebSocket subscriptions для реального часу: eth_subscribe("newHeads") дає нові блоки, eth_subscribe("logs", filter) — потокові події. Критично: при реконекті завжди робити catch-up через polling від останнього обробленого блоку.
Firehose (StreamingFast/Pinax) — бінарний протокол поверх gRPC, спеціально для high-throughput індексації. Швидкість інгестії на порядок вища за JSON-RPC. Використовується в Substreams. Якщо потрібно обробити 15M+ блоків Ethereum — розглядаємо в першу чергу.
Transform: трансформація та збагачення
Це найоб'ємніший шар за логікою. Завдання:
Декодування ABI. Raw log — це topics[] (bytes32) і data (bytes). Декодування через viem/ethers/web3.py. Нюанс з proxy-контрактами: ABI потрібно брати від implementation, а не proxy. EIP-1967 визначає стандартний slot 0x360894a13ba1a3210667c828492db98dca3e2076cc3735a920a3ca505d382bbc для адреси імплементації.
import { decodeEventLog, parseAbiItem } from 'viem' // Для proxy: резолвимо implementation const implSlot = '0x360894a13ba1a3210667c828492db98dca3e2076cc3735a920a3ca505d382bbc' const implAddr = await client.getStorageAt({ address: proxy, slot: implSlot }) const event = parseAbiItem('event Swap(address indexed sender, address indexed recipient, int256 amount0, int256 amount1, uint160 sqrtPriceX96, uint128 liquidity, int24 tick)') const decoded = decodeEventLog({ abi: [event], data: log.data, topics: log.topics }) Збагачення даних (enrichment). Чисті події рідко містять все необхідне. Типові збагачення:
- USD-вартість: підтягуємо ціну токена на момент блоку з Chainlink або власного price oracle
- Метадані токенів:
symbol(),decimals()— кешуємо агресивно, вони immutable - Identity resolution: маппінг адрес на відомі протоколи (Uniswap Router, Aave Pool)
Нормалізація. Суми токенів приводимо до decimal з правильною кількістю знаків. uint256 з контракту → Python Decimal або PostgreSQL numeric — ніяких float64, втратите точність на великих значеннях.
Трансформації зі станом — найскладніше. Обчислення running total, поточних балансів, позицій LP. Вимагає чіткого порядку обробки подій всередині блоку (сортування за logIndex).
Load: запис у сховище
Батчевий запис — обов'язково. Не INSERT по одному запису. PostgreSQL COPY або bulk INSERT через executemany:
# 10-50x швидше одиночних INSERT await conn.executemany( """ INSERT INTO swaps (block_number, tx_hash, log_index, pool, sender, amount0, amount1, price_usd, ts) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9) ON CONFLICT (tx_hash, log_index) DO NOTHING """, [(s.block, s.tx_hash, s.log_index, s.pool, s.sender, s.amount0, s.amount1, s.price, s.ts) for s in batch] ) ON CONFLICT DO NOTHING — страховка від дублів при retry після помилки. Завжди додавати UNIQUE(tx_hash, log_index).
Як правильно обробляти реорганізації блокчейну?
Реорг на Ethereum — не виняткова ситуація. На PoS-Ethereum реорги глибиною 1-2 блоки трапляються кілька разів на день. Ігнорувати — значить мати "забруднені" дані в базі.
Стратегія: tombstone + replay. Кожен запис містить block_hash. При отриманні нового блоку перевіряємо, чи не змінився block_hash для вже обробленого block_number:
-- Виявлення реоргу SELECT block_number, block_hash FROM processed_blocks WHERE block_number >= $1 AND block_hash != ANY($2::bytea[]) ORDER BY block_number; -- При розбіжності в одній транзакції: BEGIN; DELETE FROM swaps WHERE block_hash = ANY($orphaned_hashes); DELETE FROM processed_blocks WHERE block_hash = ANY($orphaned_hashes); INSERT INTO processed_blocks ...; INSERT INTO swaps ...; COMMIT; Для фінансових даних — чекаємо safe finality (12+ блоків на PoS-Ethereum) перед тим як вважати дані достовірними. Для аналітики — latest достатньо з позначкою "preliminary".
Які інструменти для черги та оркестрації обрати?
Для нетривіальних пайплайнів потрібна черга між Extract та Transform/Load — буфер при піковому навантаженні та ізоляція відмов.
| Інструмент | Коли використовувати |
|---|---|
| Redis Streams | < 10k подій/сек, проста топологія, потрібна швидкість розробки |
| Apache Kafka | > 10k подій/сек, кілька consumer groups, retention для replay |
| RabbitMQ | Складна маршрутизація, fanout на кілька downstream |
| Celery + Redis | Разові задачі, немає вимог до throughput |
Для більшості DeFi-проектів Redis Streams достатньо. Kafka додає операційної складності, але дає можливість replay — читати історію заново при додаванні нової трансформації.
Оркестрація з Airflow або Prefect потрібна коли пайплайн має залежності: спочатку завантаж ціни, потім рахуй USD-вартість свапів. DAG описує ці залежності явно.
Схема бази даних
Критичні рішення по схемі:
Партиціонування за часом — обов'язково для таблиць подій. PostgreSQL native partitioning або TimescaleDB hypertables. Без партиціонування VACUUM на таблиці з 500M рядків займе години і заблокує інсерти.
-- TimescaleDB: автоматичне партиціонування за часом SELECT create_hypertable('swaps', 'block_time', chunk_time_interval => INTERVAL '1 day'); -- Компресія старих чанків SELECT add_compression_policy('swaps', INTERVAL '7 days'); Індекси тільки потрібні. Кожен індекс — це overhead на INSERT. Типовий набір:
-
(pool_address, block_time)— запити по конкретному пулу за період -
(sender, block_time)— історія транзакцій користувача -
(tx_hash, log_index)— UNIQUE constraint для ідемпотентності
Materialized views для агрегатів. Не рахуйте суми обсягів на льоту по 100M рядків. Materialized view з daily/hourly агрегатами + REFRESH MATERIALIZED VIEW CONCURRENTLY за розкладом.
Продуктивність: реальні числа
Для орієнтиру: пайплайн на Python + asyncio + PostgreSQL на сервері 8 CPU / 32 GB RAM обробляє ~2000-5000 подій/сек при записі. Для історичної синхронізації Ethereum (2M+ блоків) це означає кілька діб роботи.
Оптимізації в порядку impact:
- Паралельна інгестія — кілька воркерів на різних діапазонах блоків. Прискорення лінійне до числа CPU та лімітів RPC.
- Вимкнення індексів при bulk load — завантажуємо сирі дані, потім
CREATE INDEX CONCURRENTLY. 3-10x прискорення вставки. - Перехід на Rust/Go для критичних компонентів. Парсинг ABI та десеріалізація блоків в Rust (
alloycrate) швидше Python в 10-20x. - Firehose замість JSON-RPC — якщо доступний для потрібної мережі, дає 5-10x прискорення інгестії.
Економія до 40% часу на історичній синхронізації за рахунок паралельної інгестії — перевірено на проектах з навантаженням 10 000 подій/сек.
Моніторинг пайплайну
Метрики які мають бути з першого дня:
- Pipeline lag —
current_block - processed_block. Алерт при > 20 блоків. Зростання lag означає вузьке місце десь у ланцюжку. - Reorg rate — кількість реорганізацій за годину. Різке зростання = нестабільна нода або RPC.
- Throughput — подій/сек на кожному етапі. Дозволяє знайти bottleneck.
- Error rate — кількість помилок декодування. > 0 означає невідомий ABI або змінений контракт.
Технологічний стек
| Компонент | Вибір | Альтернатива |
|---|---|---|
| Мова | Python (asyncio + web3.py) | TypeScript/Node.js (viem), Rust (alloy) |
| Високопродуктивна інгестія | Substreams + Firehose | Custom Rust ingester |
| Черга | Redis Streams | Apache Kafka |
| База даних | PostgreSQL 16 + TimescaleDB | ClickHouse (тільки аналітика) |
| Оркестрація | Prefect / Airflow | Temporal (складні workflows) |
| Моніторинг | Prometheus + Grafana | Datadog |
Процес розробки
Фаза 1 (3-5 днів): проектування. Визначаємо джерела даних, контракти та події, схему БД, вимоги до latency та обсягів. Прототип інгестора на тестових даних.
Фаза 2 (7-14 днів): ядро пайплайну. Extract + Transform + Load з обробкою реорганізацій. Тестування на mainnet-даних, верифікація коректності через порівняння з on-chain станом.
Фаза 3 (3-5 днів): продуктивність. Профілювання, оптимізація bottleneck-ів, налаштування БД (індекси, партиціонування, vacuum).
Фаза 4 (2-3 дні): деплой та моніторинг. Docker Compose або Kubernetes, налаштування алертів, runbook.
Разом: 2-4 тижні для пайплайну одного протоколу. Мультичейн з крос-чейн агрегацією — 4-8 тижнів. Вартість розробки розраховується індивідуально. Аудит існуючого пайплайну — за запитом. Отримайте консультацію по вашому проекту — зв'яжіться з нами для безкоштовної оцінки.
Що входить у роботу
- Документація архітектури та схеми даних
- Вихідний код пайплайну (GitHub)
- Налаштований моніторинг (Grafana дашборди, Prometheus алерти)
- Інструкція з деплою та експлуатації
- Навчання команди (2-3 сесії)
- Підтримка протягом 1 місяця після запуску
Наша команда має 8+ років досвіду в блокчейн-розробці та більше 50 завершених ETL-проектів. Ми використовуємо лише перевірені інструменти та гарантуємо надійність пайплайну навіть при пікових навантаженнях. Замовте розробку або аудит — отримайте готове рішення в стислі терміни.







