Надійний 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 необхідно:

  1. Десеріалізувати транзакції — розкодувати CompactArray інструкцій, витягнути дані з Anchor-програм (якщо є discriminator), розпарсити PDA-адреси.
  2. Нормалізувати формати — привести час до UTC (slot → приблизний timestamp через банк-хеш або зовнішній оракул), уніфікувати адреси до base58.
  3. Збагатити контекстом — додати ім'я програми за відомим 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 містив блоки за даними експлорера.

Механізм обробки:

  1. Відстеження останнього обробленого слота — зберігайте у окремій таблиці (pipeline_state) пару (last_slot, last_blockhash).
  2. Виявлення прогалини — якщо різниця між поточним слотом і last_slot перевищує поріг (наприклад, 5 слотів), ініціюйте перевірку.
  3. Заповнення через getBlock — викличте RPC-метод getBlock для кожного пропущеного слота. Якщо блок існує — обробіть його. Якщо слот порожній — оновіть last_slot без запису даних.
  4. Обмеження ретроспективи — не намагайтеся заповнювати прогалини глибше ніж на 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-з'єднання (валідатор перезавантажився, мережевий розрив):

  1. Pipeline автоматично переходить у режим polling через RPC з інтервалом 400 мс.
  2. Після відновлення Geyser-з'єднання pipeline порівнює last_slot з поточним слотом і заповнює прогалину.
  3. Якщо прогалина перевищує 1000 слотів — алерт оператору, ручна перевірка необхідності повного ресинкронізування.

Збій цільової бази даних (PostgreSQL failover, втрата з'єднання):

  1. Transform-шар продовжує отримувати дані та буферизує їх у локальну чергу (in-memory з лімітом, наприклад, 50 000 записів).
  2. При переповненні буфера — найстаріші записи скидаються з логуванням, а last_slot не оновлюється. Це навмисний вибір: краще втратити буфер і відновитися через прогалину, ніж вивести з ладу весь pipeline через OOM.
  3. Після відновлення бази — спочатку записується буфер, потім заповнюються прогалини.

Логічна помилка в transform (неправильний парсинг нової версії програми):

  1. Error rate зростає, алерт спрацьовує.
  2. Оператор призупиняє pipeline (graceful stop: завершити обробку поточного блоку, зберегти стан).
  3. Виправляється логіка парсингу.
  4. Pipeline відкочується до останнього коректного слоту (з pipeline_state) і перезапускається.
  5. Уражені записи в цільовій базі видаляються за діапазоном слотів і перезаписуються.

Повний ресинкронізування

Якщо прогалина занадто велика або дані пошкоджено, єдиний надійний шлях — повне ресинкронізування з певного слоту:

  1. Зупиніть pipeline та збережіть поточний last_slot.
  2. Видаліть дані з цільової таблиці за діапазоном слотів (DELETE WHERE slot >= target_slot). Для великих таблиць використовуйте парціальне видалення партіями, щоб не заблокувати таблицю.
  3. Скиньте стан — оновіть pipeline_state до target_slot.
  4. Запустіть pipeline у режимі backfill — той самий код, але з пріоритетом швидкості над реальним часом. У цьому режимі можна збільшити розмір партій для запису та використовувати паралельні запити до RPC.
  5. Переключіть у 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 режим. Виявлена розбіжність означає, що відновлення неповне і треба розширити діапазон ресинкронізування.

Джерела