אנו מפתחים data lake לבלוקצ'יין — שכבה שפותרת את הבעיה המרכזית של אחסון נתוני on-chain: נתוני בלוקצ'יין גולמיים אינם מתאימים לשאילתות מורכבות. צמתי JSON-RPC עונים על "מה קרה בבלוק X" אך לא על "הראה את כל ההמרות של Uniswap V3 ב-30 הימים האחרונים עבור כתובות עם נפח גבוה". ה-data lake הופך בלוקים גולמיים לטבלאות מובנות, מתויגות וניתנות לשאילתה מהירה.
ברשת Ethereum הראשית כיום — כ-20 מיליון בלוקים, בערך 2 מיליארד עסקאות, וטרה-בייטים של לוגי אירועים. ההיסטוריה המלאה של Ethereum בפורמט Parquet תופסת 3–4 TB. בכל בלוק חדש (כל 12 שניות), מתווספות מאות עסקאות ואלפי רשומות לוג. הוסיפו את BSC, Polygon, Arbitrum, Base — לכל רשת יש היסטוריה וקצב גידול משלה. ה-data lake שלנו מאחד אותן לסביבה אנליטית אחת.
למה יש צורך ב-Data Lake לניתוח בלוקצ'יין
שלוש מחלקות נתונים עם מאפיינים שונים:
בלוקים ועסקאות — סכמה מובנית וצפויה. האתגר המרכזי: reorgs — מזלגות זמניים שלאחריהם השרשרת נכתבת מחדש. צינור הנתונים של EVM חייב להיות מסוגל להחזיר לאחור נתונים שכבר נכתבו.
לוגי אירועים — החשובים ביותר לאנליטיקה. Transfer, Swap, Liquidation, Mint — כל אירועי ה-EVM. בעיה: פענוח ABI. ללא ABI של החוזה, לוג הוא רק בייטים עם topics. אנו יוצרים רישום ABI דרך Etherscan API ו-Sourcify כדי לפענח מיליוני אירועים אוטומטית.
Traces (עסקאות פנימיות) — קריאות בין חוזים שאינן יוצרות עסקה ישירה. ללא traces, חלק משמעותי מ-DeFi אינו נראה: הלוואות flash בתוך עסקה אחת, liquidations רקורסיביות, חבילות MEV. קבלת traces דרך debug_traceTransaction היא פעולה כבדה, הזמינה רק בצמתי archive.
טיפול ב-Reorg
Reorg הוא כאב הראש המרכזי של כל צינור אינדוקס בלוקצ'יין. Ethereum עם Proof-of-Stake יש לו סופיות הסתברותית לאחר כמה בלוקים וסופיות מלאה לאחר כ-12.8 דקות (2 epochs). לרשתות L2 יש מודל מורכב עוד יותר.
גישה סטנדרטית:
- כתיבת בלוקים עם השהיית אישור (המתנה ל-N אישורים לפני כתיבה לשכבה הסופית). עבור Ethereum: 32–64 בלוקים.
- שמירת שכבת staging עבור M הבלוקים האחרונים — נתונים נכתבים מיד אך מסומנים כ-
pending. - הרשמה לאירועי
Reorganizationמהצומת (WebSocketnewHeads+ השוואת parentHash). בעת reorg — מחיקת בלוקים מושפעים מ-staging והחלת השרשרת החדשה מחדש.
עבור Iceberg זה נפתר באלגנטיות באמצעות time travel ופעולות merge. עבור ClickHouse — באמצעות ReplacingMergeTree עם עמודת גרסה.
ארכיטקטורת Data Lake
שכבת קליטה (Ingestion)
שתי גישות:
- קליטה מבוססת צומת — חיבור ישיר לצומת דרך WebSocket. הרשמה לבלוקים חדשים + backfill דרך קריאות אצווה
eth_getLogs. דורש צומת archive. למילוי מיליוני בלוקים, אנו משתמשים בעיבוד מקבילי עם asyncio. - ספקי נתונים חיצוניים — Goldsky, Envio, Substreams. התחלה מהירה יותר אך נעילת ספק ועלות גבוהה יותר בקנה מידה.
אחסון: בחירת פורמט ומנוע
עבור נתוני בלוקצ'יין גולמיים, אחסון עמודי (columnar) הוא אופטימלי:
- Apache Parquet על S3/GCS — התקן. דחיסת zstd מקטינה את הנפח פי 5–10. חלוקה לפי תאריך ומספר בלוק.
- Apache Iceberg על גבי Parquet — ACID, אבולוציית סכמה, time travel. קריטי עבור reorgs.
- ClickHouse — OLAP לשאילתות חמות. מאות מיליוני שורות בשניות. שימוש ב-ClickHouse לנתוני בלוקצ'יין מהיר עד פי 1000 משאילתת צמתי JSON-RPC ישירות.
ארכיטקטורה טיפוסית דו-שכבתית:
Raw layer (S3 + Parquet/Iceberg) ↓ ETL (dbt / Spark / Flink) Serving layer (ClickHouse / BigQuery) ↓ Query API Analytics / Trading systems / Dashboards פענוח ABI והעשרה
לוגי אירועים גולמיים מכילים topics (חתימות אירועים) ונתונים (מקודדים ב-ABI). לפענוח אירועי Ethereum, יש צורך ברישום 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 — טבלה הממפה Raw layer (S3 + Parquet/Iceberg) ↓ ETL (dbt / Spark / Flink) Serving layer (ClickHouse / BigQuery) ↓ Query API Analytics / Trading systems / Dashboards . מקורות: Etherscan API, Sourcify, 4byte.directory. חוזים לא ידועים מעובדים כברייטים גולמיים, והעשרה מתבצעת כשזמין ABI.
העשרת מטא-דאטה של טוקנים: עבור העברות ERC-20 אנו זקוקים ל-decimals, symbol, price. מחירים נלקחים מרשומות TWAP של Uniswap V3 או מ-APIs חיצוניים (נתונים היסטוריים).
סכמת נתונים וטבלאות מפתח
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); טבלאות נפרדות לסוגי אירועים בתדירות גבוהה: 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} , contract_address → ABI, 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); . חלוקה לפי חודש.
דוגמה לחלוקה ואופטימיזציה
עבור רשת Ethereum, אנו מחלקים את טבלת האירועים לפי חודש. זה מאפשר מחיקה מהירה של נתונים מיושנים וסריקה יעילה של טווחי זמן. אינדקסים של Bloom filter על החוזה מאיצים סינון לפי כתובת.שאילתה לדוגמה:
erc20_transfers מתבצעת בתוך פחות משנייה על 50M שורות.
השוואת שיטות רכישת נתונים
| פרמטר | מבוסס צומת | צד שלישי (Goldsky) |
|---|---|---|
| מהירות התחלה | בינונית (הקמת צומת) | גבוהה (מפתח API) |
| שליטה בנתונים | מלאה | מוגבלת על ידי הספק |
| עלות בקנה מידה | נמוכה (צמתים עצמיים) | גבוהה (תשלום לפי נפח) |
| טיפול ב-Reorg | מנגנון מותאם אישית | מובנה (אך אטום) |
מה כלול בעבודה
אתם מקבלים:
- Data lake עובד עם הרשתות והאירועים שנבחרו.
- תיעוד של סכמת הנתונים וצינור ה-ETL.
- גישה ל-ClickHouse (או שכבת שירות אחרת) עם דוגמאות שאילתות.
- ניטור השהיה והתראות לבעיות.
- הכשרת צוות וחודש תמיכה.
העלות הטיפוסית של פרויקט נעה בין $15,000 ל-$50,000, תלוי במספר הרשתות ומורכבות האירועים. יכולת הרחבה: הוספת חוזים או רשתות חדשות נעשית דרך קונפיגורציה ללא שינויי קוד.
שלבי פיתוח
| שלב | תוכן | משך |
|---|---|---|
| עיצוב | הגדרת היקף, רשתות/אירועים, סכמת נתונים | 1–2 שבועות |
| קליטת ליבה | מאזין WebSocket, backfill, מטפל ב-reorg | 3–4 שבועות |
| רישום ABI | צבירת ABI, פענוח, העשרה | 2–3 שבועות |
| שכבת אחסון | Parquet/Iceberg, ClickHouse, ETL | 3–4 שבועות |
| API שירות | REST/GraphQL, הגבלת קצב | 2–3 שבועות |
| ניטור ותפעול | Airflow, התראות, תיעוד | 1–2 שבועות |
למה לעבוד איתנו
המומחיות שלנו — מעל 5 שנים בהנדסת בלוקצ'יין, יותר מ-20 צינורות נתונים שיושמו עבור פרוטוקולי DeFi וקרנות קריפטו. אנו מפתחים מקצה לקצה — מעיצוב סכמה ועד פריסה וניטור. נבחן את הפרויקט שלכם תוך יומיים — צרו קשר לייעוץ. הזמינו פיתוח data lake מקצה לקצה כדי להאיץ אנליטיקת נתוני on-chain.
מילות אמון: אחריות לאיכות, מומחים מוסמכים, ניסיון רב ב-L1/L2. מידע נוסף על מבנה הבלוקצ'יין ניתן לקרוא ב-ויקיפדיה: בלוקצ'יין.







