Надійний data pipeline для Solana — це система, яка безперервно отримує блоки й транзакції, нормалізує їх і записує в цільове сховище з гарантованою цілісністю та можливістю відновлення після збою. Нижче — покрокова архітектура та інженерні рішення, які доведено в production-середовищах.
Архітектура pipeline: ingestion, transform, load
Класичний ETL-підхід для Solana потребує адаптації через специфіку блокчейна: фінальність після підтвердження, форки, паралельне виконання транзакцій та обмеження RPC-нод.
Ingestion: отримання даних
Є три перевірені способи ingestion, кожен із своїми межами:
- Geyser Plugin — підключається безпосередньо до валідатора через плагінний інтерфейс. Отримує повні блоки, транзакції, облікові записи (accounts) у реальному часі без RPC-обмежень. Вимагає власну або орендовану ноду з увімкненим geyser-інтерфейсом. Це найнадійніший варіант для high-throughput pipeline.
- WebSocket-підписки (onLogs, onAccountChange) — підходять для вузьких завдань: моніторинг конкретних програм або гаманців. Мають обмеження за кількістю підписок на одному з'єднанні та не гарантують доставку всіх подій за високого навантаження.
- Пакетні RPC-запити — для історичних даних або періодичного синхронізування. Детальніше про цей підхід описано в окремому матеріалі про RPC 2.0 для пакетних запитів.
Для production-контексту рекомендована конфігурація: Geyser Plugin як основне джерело, пакетний RPC як резервний механізм для заповнення прогалин.
Transform: нормалізація та збагачення
Сира дата з Geyser — це protobuf- або JSON-структури, які містять serialized transaction data, метадані блоку та логи інструкцій. На етапі transform необхідно:
- Десеріалізувати транзакції — розкодувати CompactArray інструкцій, витягнути дані з Anchor-програм (якщо є discriminator), розпарсити PDA-адреси.
- Нормалізувати формати — привести час до UTC (slot → приблизний timestamp через банк-хеш або зовнішній оракул), уніфікувати адреси до base58.
- Збагатити контекстом — додати ім'я програми за відомим program ID, розкласти токен-трансфери з інструкцій Token Program або SPL Transfer Hook, прив'язати транзакцію до відповідного блоку.
Критичне правило: transform-шар має бути безстаним (stateless) щодо бізнес-логіки. Кожен блок обробляється незалежно, а стан формується лише на етапі load.
Load: запис у цільове сховище
Вибір сховища визначається шаблоном доступу:
| Сховище | Підходить для | Обмеження |
|---|---|---|
| PostgreSQL з partitioning за slot | Аналітичні запити, JOIN по транзакціях і облікових записах | Обмежена пропускна здатність на запис при високому піковому навантаженні |
| ClickHouse | Часові ряди, агрегації, логи інструкцій | Не підходить для транзакційних оновлень (UPDATE/DELETE) |
| Apache Kafka / Redpanda | Буфер між transform та споживачами, event sourcing | Потребує додаткової інфраструктури та налаштування ретенції |
| TimescaleDB | Метрики валідатора, часові ряди облікових записів | Спеціалізоване, не замінює повноцінну OLTP-базу |
Для надійного запису використовуйте batch insert з розміром партії 500–2000 записів та інтервалом флашу 1–3 секунди. Менші партії створюють надмірне навантаження на базу, більші — збільшують ризик втрати даних при збої.
Обробка пропусків та дублікатів
Solana-пайплайн стикається з двома типами проблем цілісності: прогалини в послідовності слотів та дублікати через форки.
Пропуски слотів
Не кожен слот містить блоки (пропущені лідерські слоти — нормальне явище). Пропуск — це ситуація, коли pipeline очікує слот N, а отримує N+2, причому N+1 містив блоки за даними експлорера.
Механізм обробки:
- Відстеження останнього обробленого слота — зберігайте у окремій таблиці (pipeline_state) пару (last_slot, last_blockhash).
- Виявлення прогалини — якщо різниця між поточним слотом і last_slot перевищує поріг (наприклад, 5 слотів), ініціюйте перевірку.
- Заповнення через getBlock — викличте RPC-метод getBlock для кожного пропущеного слота. Якщо блок існує — обробіть його. Якщо слот порожній — оновіть last_slot без запису даних.
- Обмеження ретроспективи — не намагайтеся заповнювати прогалини глибше ніж на 500–1000 слотів назад. За цією межею дані можуть бути недоступні на стандартних RPC-нодах.
Дублікати через форки
Solana використовує оптимістичне підтвердження: блок може бути замінений форком протягом кількох слотів. Pipeline отримує обидва варіанти блоку.
Стратегія deduplication:
- На рівні запису — використовуйте унікальний constraint по (slot, blockhash) у цільовій таблиці. При вставці дубліката база відхилить запис без помилки (INSERT ... ON CONFLICT DO NOTHING).
- На рівні бізнес-логіки — якщо обробляєте транзакції, використовуйте signature як унікальний ключ. Одна транзакція може потрапити в різні форки, але signature залишається незмінним.
- Форкове очищення — після досягнення фінальності (за замовчуванням 32 підтвердження, але перевірте актуальне значення в документації) можна запустити фоновий процес, який видаляє блоки з нефінальними гілками. Для більшості аналітичних завдань це не обов'язково — достатньо ігнорувати їх у запитах за прапорцем is_finalized.
Моніторинг та алерти для pipeline
Моніторинг pipeline відрізняється від моніторингу валідатора: ви відстежуєте не здоров'я ноди, а цілісність та своєчасність потоку даних.
Ключові метрики
- Slot lag — різниця між найвищим відомим слотом (з RPC: getSlot) та останнім обробленим слотом у pipeline. Нормальне значення: 0–3 слоти. Зростання понад 10 слотів означає, що pipeline не встигає.
- Throughput (транзакцій/с) — кількість оброблених транзакцій за секунду. Порівнюйте з фактичною пропускною здатністю мережі (можна отримати через getRecentPerformanceSamples).
- Error rate — частота помилок десеріалізації, невдалих вставок у базу, таймаутів RPC. Окремо відстежуйте помилки за типами: parse errors вказують на зміну формату даних, connection errors — на інфраструктурні проблеми.
- Gap count — кількість виявлених прогалин за останню хвилину. Більше нуля — привід для перевірки.
- Consumer lag — якщо використовуєте чергу (Kafka), відстежуйте відставання споживачів від продюсера.
Пороги алертів
| Метрика | Warning | Critical |
|---|---|---|
| Slot lag | > 5 слотів | > 20 слотів |
| Error rate | > 1% за 5 хвилин | > 5% за 5 хвилин |
| Gap count | > 0 за хвилину | > 3 за хвилину без зменшення |
| Throughput drop | < 80% від очікуваного | < 50% від очікуваного |
Алерти налаштовуйте через Prometheus + Alertmanager або еквівалентний стек. Уникайте алертів, які спрацьовують на кожен пропущений слот — це створить шум, оскільки порожні слоти є нормою.
Дашборд
Мінімальний дашборд для оператора pipeline має містити: графік slot lag у часі, throughput з порівнянням з мережею, кількість активних прогалин, error rate за типами, стан з'єднання з джерелом даних (Geyser/RPC). Усі графіки — з інтервалом агрегації не більше 15 секунд для виявлення різких спадів.
План відновлення після збою
Збій pipeline — це не питання «чи станеться», а «коли». План відновлення має бути задокументований і перевірений до першого реального інциденту.
Типи збоїв та відповіді
Збій Geyser-з'єднання (валідатор перезавантажився, мережевий розрив):
- Pipeline автоматично переходить у режим polling через RPC з інтервалом 400 мс.
- Після відновлення Geyser-з'єднання pipeline порівнює last_slot з поточним слотом і заповнює прогалину.
- Якщо прогалина перевищує 1000 слотів — алерт оператору, ручна перевірка необхідності повного ресинкронізування.
Збій цільової бази даних (PostgreSQL failover, втрата з'єднання):
- Transform-шар продовжує отримувати дані та буферизує їх у локальну чергу (in-memory з лімітом, наприклад, 50 000 записів).
- При переповненні буфера — найстаріші записи скидаються з логуванням, а last_slot не оновлюється. Це навмисний вибір: краще втратити буфер і відновитися через прогалину, ніж вивести з ладу весь pipeline через OOM.
- Після відновлення бази — спочатку записується буфер, потім заповнюються прогалини.
Логічна помилка в transform (неправильний парсинг нової версії програми):
- Error rate зростає, алерт спрацьовує.
- Оператор призупиняє pipeline (graceful stop: завершити обробку поточного блоку, зберегти стан).
- Виправляється логіка парсингу.
- Pipeline відкочується до останнього коректного слоту (з pipeline_state) і перезапускається.
- Уражені записи в цільовій базі видаляються за діапазоном слотів і перезаписуються.
Повний ресинкронізування
Якщо прогалина занадто велика або дані пошкоджено, єдиний надійний шлях — повне ресинкронізування з певного слоту:
- Зупиніть pipeline та збережіть поточний last_slot.
- Видаліть дані з цільової таблиці за діапазоном слотів (DELETE WHERE slot >= target_slot). Для великих таблиць використовуйте парціальне видалення партіями, щоб не заблокувати таблицю.
- Скиньте стан — оновіть pipeline_state до target_slot.
- Запустіть pipeline у режимі backfill — той самий код, але з пріоритетом швидкості над реальним часом. У цьому режимі можна збільшити розмір партій для запису та використовувати паралельні запити до RPC.
- Переключіть у real-time — коли slot lag досягне нуля, переведіть pipeline у стандартний режим.
Час повного ресинкронізування залежить від глибини відкату та пропускної здатності RPC. Орієнтовно: 100 000 слотів через пакетні запити займають від 20 до 60 хвилин за умови доступного RPC з високими лімітами. Перевірте актуальні ліміти вашого RPC-провайдера перед розрахунком часу відновлення.
Перевірка цілісності після відновлення
Після будь-якого відновлення обов'язково виконайте верифікацію:
- Порівняйте кількість транзакцій за діапазоном слотів у вашій базі з даними експлорера (наприклад, через публічний API).
- Перевірте відсутність розривів у послідовності слотів: SELECT slot + 1 FROM blocks EXCEPT SELECT slot FROM blocks — має повернути порожній результат для обробленого діапазону.
- Перевірте унікальність signature у межах діапазону відновлення.
Ці перевірки інтегруйте в процес відновлення як автоматичний крок перед переключенням у real-time режим. Виявлена розбіжність означає, що відновлення неповне і треба розширити діапазон ресинкронізування.