צינור ETL על-רשת: בנייה לנתוני בלוקצ'יין

ביום השלישי לאחר השקת האנליטיקה של Uniswap v3, אתה מגלה ש-`eth_getLogs` עם פילטר רחב מתחיל לפג תוקף, צבירות סוטות עקב ארגונים מחדש שהוחמצו, וה-PostgreSQL שלך מתנפח עם טבלאות ג'יגה-בייט ללא חלוקה. צינור ETL על-רשת אינו רק "קריאת לוגים ועיבוד

שירותי פיתוח בלוקצ'יין

שאלות נפוצות

העבודות האחרונות

  • image_website-b2b-advance_0.webp
    פיתוח אתר חברה B2B ADVANCE
    1452
  • image_web-applications_feedme_466_0.webp
    פיתוח אפליקציית ווב עבור FEEDME
    1310
  • image_websites_belfingroup_462_0.webp
    פיתוח אתר עבור BELFINGROUP
    1005
  • image_ecommerce_furnoro_435_0.webp
    פיתוח חנות מקוונת לחברת FURNORO
    1270
  • image_logo-advance_0.webp
    עיצוב לוגו לחברת B2B Advance
    719
  • image_crm_enviok_479_0.webp
    פיתוח אפליקציית ווב עבור Enviok
    1012

ביום השלישי לאחר השקת ניתוחי Uniswap v3, אתה מגלה ש-eth_getLogs עם פילטר רחב מתחיל לתפוס פסק זמן, האגרגציות מתפצלות עקב ארגונים מחדש שהוחמצו, וה-PostgreSQL שלך מתנפח עם טבלאות בג'יגה-בייט ללא חלוקה למחיצות. צינור ETL על-השרשרת אינו רק "קריאת לוגים וכתיבה למסד נתונים". זו מערכת עם ערבויות עקביות, טיפול בארגון מחדש, טרנספורמציה של נתונים וניהול עומסים. אנחנו בונים את זה נכון מהפעם הראשונה, ובטקסט הזה ננתח את ההחלטות הארכיטקטוניות המרכזיות.

דוגמה מהפרקטיקה: אחד הלקוחות שלנו איבד שבועיים על שחזור נתונים עקב טיפול שגוי בארגון מחדש. לאחר יישום הצינור שלנו, חיסכון בזמן הסתכם ב-40% בסנכרון היסטורי, ועלות הבעלות על התשתית ירדה ב-30% (5,000 דולר לחודש בממוצע) עקב אופטימיזציה של אחסון.

ETL ללא טיפול בארגון מחדש אינו ETL, אלא ייצור זבל. הניסיון שלנו אומר: 90% מהבעיות נפתרות על ידי ארכיטקטורה נכונה בהתחלה.

למה צריך צינור ETL על-השרשרת

צינור ETL על-השרשרת מחלץ נתונים גולמיים מהבלוקצ'יין (לוגי אירועים, טרנזקציות פנימיות, שינויי מצב), הופך אותם לרשומות מובנות (פענוח ABI, העשרת מחירים, נורמליזציה של כמויות) וטוען אותם לאחסון אנליטי. ללא צינור כזה, אי אפשר לבנות לוחות מחוונים לפרוטוקולי DeFi, לעקוב אחר נזילות בזמן אמת או לבצע ניתוח היסטורי. אתגרים עיקריים: ארגונים מחדש של השרשרת, נפחים עצומים (עד 15M+ בלוקים באת'ריום, 500 GB+ של לוגים גולמיים), והצורך להבטיח עקביות במהלך קליטה מקבילית.

איך הארכיטקטורה עובדת: שלוש שכבות ETL

ETL קלאסי (חילוץ — טרנספורמציה — טעינה) בהקשר הבלוקצ'יין מקבל מאפיינים מיוחדים: מקור הנתונים הוא בלתי-משתנה אך לא סופי (ארגונים מחדש), הנפחים נמדדים במאות מיליוני אירועים, והשהייה יכולה לנוע בין שניות לשעות בהתאם למשימה.

חילוץ: קליטה מהצומת

בחירת מקור הנתונים קובעת את כל השאר. שלוש רמות עם מורכבות גוברת:

  • לוגים/אירועים — מה שהחוזה פולט במפורש. זול, מהיר, מובנה דרך ABI. מגבלה: רק מה שהמפתח בחר לתעד.
  • טרייסים (טרנזקציות פנימיות) — כל הקריאות בתוך טרנזקציה, כולל העברות ETH ללא אירועים. דורש debug_traceTransaction או trace_block (בסגנון Parity). לא כל הצמתים תומכים בזה; Erigon היא הבחירה הטובה ביותר למשימות כבדות בטרייסים.
  • הבדלי מצב — שינויים בחריצי אחסון לכל בלוק. שלמות מקסימלית, אבל נפח נתונים עצום וקושי בפירוש ללא ABI.

לרוב משימות ה-DeFi, לוגים + טרייסים מספיקים. הבדלי מצב נחוצים לניתוחי MEV ולניטור חוזים ללא אירועים (למשל, 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 לזמן אמת: # 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 נותן בלוקים חדשים, eth_subscribe("newHeads") — אירועים בזרימה. קריטי: בחיבור מחדש, תמיד בצע השלמה באמצעות פולינג מהבלוק האחרון שעובד.

Firehose (StreamingFast/Pinax) — פרוטוקול בינארי על גבי gRPC, במיוחד לאינדוקס בתפוקה גבוהה. מהירות קליטה בסדר גודל גבוהה יותר מ-JSON-RPC. משמש ב-Substreams. אם אתה צריך לעבד 2M+ בלוקים של את'ריום, שקול את זה קודם.

טרנספורמציה: המרה והעשרה

זו השכבה הגדולה ביותר בלוגיקה. משימות:

פענוח ABI. לוג גולמי מורכב מ-eth_subscribe("logs", filter) (bytes32) ו-topics[] (bytes). פענוח דרך viem/ethers/web3.py. אזהרה עם חוזי פרוקסי: יש לקחת ABI מה-implementation, לא מהפרוקסי. EIP-1967 מגדיר את החריץ הסטנדרטי data לכתובת ה-implementation.

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 }) 

העשרת נתונים. אירועים גולמיים כמעט ולא מכילים את כל מה שצריך. העשרות אופייניות:

  • ערך בדולרים: משוך מחיר טוקן בזמן הבלוק מ-Chainlink או מאורקל מחירים מותאם
  • מטא-דאטה של טוקן: 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 }) — שמור במטמון באגרסיביות, הם בלתי-משתנים
  • זיהוי זהויות: מיפוי כתובות לפרוטוקולים ידועים (Uniswap Router, Aave Pool)

נורמליזציה. כמויות טוקן מומרות ל-symbol() עם מספר הנקודות העשרוניות הנכון. decimals() מהחוזה → decimal של Python או uint256 של PostgreSQL — לעולם לא Decimal, תאבד דיוק על ערכים גדולים (למשל, $1,000,000,000,000).

טרנספורמציות מצביות — החלק הקשה ביותר. חישוב סכומים רצים, יתרות נוכחיות, עמדות LP. דורש סדר ברור של עיבוד אירועים בתוך בלוק (מיון לפי numeric).

טעינה: כתיבה לאחסון

כתיבות אצווה — חובה. לא INSERT אחד אחד. float64 של PostgreSQL או INSERT בכמות גדולה דרך logIndex:

# 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] ) 

COPY — הגנה מפני כפילויות בנסיון חוזר לאחר שגיאה. תמיד הוסף executemany.

איך לטפל נכון בארגונים מחדש של הבלוקצ'יין?

ארגון מחדש באת'ריום אינו מצב חריג. ב-PoS-Ethereum, ארגונים מחדש בעומק 1-2 בלוקים מתרחשים מספר פעמים ביום. התעלמות מהם משמעותה נתונים "מזוהמים" במסד הנתונים.

אסטרטגיה: tombstone + replay. כל רשומה מכילה # 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 עבור UNIQUE(tx_hash, log_index) שכבר עובד השתנה:

-- Обнаружение реорга 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; 

לנתונים פיננסיים, המתן ל-block_hash סופיות (12+ בלוקים ב-PoS-Ethereum) לפני שתחשיב נתונים כאמינים. לאנליטיקה, block_hash מספיק עם תווית "מקדמי".

איזה תור וכלי אורקסטרציה לבחור?

לצינורות לא טריוויאליים, יש צורך בתור בין חילוץ לטרנספורמציה/טעינה — חיץ בעומס שיא ובידוד תקלות.

כלי מתי להשתמש
Redis Streams < 10k אירועים/שנייה, טופולוגיה פשוטה, פיתוח מהיר
Apache Kafka > 10k אירועים/שנייה, קבוצות צרכנים מרובות, שמירה ל-replay
RabbitMQ ניתוב מורכב, fanout למספר יעדים במורד הזרם
Celery + Redis משימות חד-פעמיות, ללא דרישות תפוקה

לרוב פרויקטי ה-DeFi, Redis Streams מספיק. Kafka מוסיף מורכבות תפעולית אבל מאפשר replay — קריאה חוזרת של היסטוריה בעת הוספת טרנספורמציה חדשה.

אורקסטרציה עם Airflow או Prefect נחוצה כאשר לצינור יש תלויות: קודם טען מחירים, אחר כך חשב ערך בדולרים של החלפות. DAG מתאר את התלויות האלה במפורש.

סכמת מסד נתונים

החלטות סכמה קריטיות:

חלוקה למחיצות לפי זמן היא חובה לטבלאות אירועים. חלוקה מקורית של PostgreSQL או hypertables של TimescaleDB. ללא חלוקה, block_number על טבלה עם 500M שורות ייקח שעות ויחסום INSERT.

-- TimescaleDB: автоматическое партиционирование по времени SELECT create_hypertable('swaps', 'block_time', chunk_time_interval => INTERVAL '1 day'); -- Компрессия старых чанков SELECT add_compression_policy('swaps', INTERVAL '7 days'); 

אינדקסים רק לפי צורך. כל אינדקס הוא עומס על INSERT. סט אופייני:

  • -- Обнаружение реорга 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 — היסטוריית טרנזקציות של משתמש
  • latest — אילוץ UNIQUE לאידמפוטנטיות

תצוגות חומריות לאגרגציות. אל תחשב סכומי נפח תוך כדי תנועה על 100M שורות. תצוגה חומרית עם אגרגציות יומיות/שעתיות + VACUUM בלוח זמנים.

ביצועים: מספרים אמיתיים

להשוואה: צינור על Python + asyncio + PostgreSQL על שרת 8 CPU / 32 GB RAM מעבד ~2000-5000 אירועים/שנייה במהלך כתיבות. לסנכרון היסטורי של את'ריום (2M+ בלוקים), זה אומר מספר ימים של פעולה.

אופטימיזציות בסדר השפעה:

  1. קליטה מקבילית — מספר עובדים על טווחי בלוקים שונים. האצה ליניארית עד למספר ה-CPU ומגבלות RPC.
  2. השבת אינדקסים במהלך טעינת אצווה — טען נתונים גולמיים, ואז -- TimescaleDB: автоматическое партиционирование по времени SELECT create_hypertable('swaps', 'block_time', chunk_time_interval => INTERVAL '1 day'); -- Компрессия старых чанков SELECT add_compression_policy('swaps', INTERVAL '7 days'); . האצת INSERT פי 3-10.
  3. מעבר ל-Rust/Go לרכיבים קריטיים. פענוח ABI ודה-סריאליזציה של בלוקים ב-Rust (קרייט (pool_address, block_time)) מהירים פי 10-20 מ-Python.
  4. Firehose במקום JSON-RPC — אם זמין לרשת היעד, נותן האצת קליטה פי 5-10.

חיסכון בזמן של עד 40% בסנכרון היסטורי עקב קליטה מקבילית — הוכח על פרויקטים עם עומס של 10,000 אירועים/שנייה. עלויות תשתית יורדות ב-5,000 דולר לחודש ללקוחות שמעבדים 100+ חוזים.

ניטור צינור

מדדים שחייבים להיות מהיום הראשון:

  • פיגור צינור — (sender, block_time). התראה ב-> 20 בלוקים. פיגור גדל מצביע על צוואר בקבוק איפשהו בשרשרת.
  • שיעור ארגונים מחדש — מספר ארגונים מחדש לשעה. עלייה חדה = צומת לא יציב או RPC.
  • תפוקה — אירועים/שנייה בכל שלב. מאפשר זיהוי צווארי בקבוק.
  • שיעור שגיאות — מספר שגיאות פענוח. > 0 פירושו ABI לא ידוע או חוזה שהשתנה.

מחסנית טכנולוגית

רכיב בחירה חלופה
שפה Python (asyncio + web3.py) TypeScript/Node.js (viem), Rust (alloy)
קליטה בביצועים גבוהים Substreams + Firehose מנגנון קליטה מותאם ב-Rust
תור Redis Streams Apache Kafka
מסד נתונים PostgreSQL 16 + TimescaleDB ClickHouse (אנליטיקה בלבד)
אורקסטרציה Prefect / Airflow Temporal (זרימות עבודה מורכבות)
ניטור Prometheus + Grafana Datadog

תהליך פיתוח

שלב 1 (3-5 ימים): תכנון. קבע מקורות נתונים, חוזים ואירועים, סכמת מסד נתונים, דרישות השהייה ונפח. אב טיפוס קליטה על נתוני בדיקה.

שלב 2 (7-14 ימים): ליבת צינור. חילוץ + טרנספורמציה + טעינה עם טיפול בארגון מחדש. בדיקות על נתוני mainnet, אימות נכונות על ידי השוואה למצב על-השרשרת.

שלב 3 (3-5 ימים): ביצועים. פרופיילינג, אופטימיזציה של צווארי בקבוק, כוונון מסד נתונים (אינדקסים, חלוקה למחיצות, vacuum).

שלב 4 (2-3 ימים): פריסה וניטור. Docker Compose או Kubernetes, הגדרת התראות, runbook.

סה"כ: 2-4 שבועות לצינור פרוטוקול יחיד. Multi-chain עם אגרגציה חוצת-שרשרת — 4-8 שבועות. עלות הפיתוח מחושבת באופן אישי, בדרך כלל 15,000-30,000 דולר לפרוטוקול יחיד. ביקורת על צינור קיים — לפי בקשה, החל מ-5,000 דולר. קבל ייעוץ לפרויקט שלך — צור קשר להערכה חינם.

מה כלול

  • תיעוד ארכיטקטורה וסכמת נתונים
  • קוד מקור של הצינור (GitHub)
  • ניטור מוגדר (לוחות מחוונים של Grafana, התראות Prometheus)
  • הוראות פריסה ותפעול
  • הדרכת צוות (2-3 מפגשים)
  • תמיכה לחודש אחד לאחר ההשקה

לצוות שלנו יש 8+ שנות ניסיון בפיתוח בלוקצ'יין ויותר מ-50 פרויקטי ETL שהושלמו. אנחנו משתמשים רק בכלים מוכחים ומבטיחים אמינות צינור גם בעומסי שיא. הזמן פיתוח או ביקורת — קבל פתרון מוכן בזמן קצר.