איסוף נתוני מסחר בזמן אמת מבורסות קריפטו
לקוח איבד לאחרונה שלושה ימים בניסיון לאסוף עסקאות מ-Binance באמצעות REST—הוא פספס 15% מהעסקאות עקב מגבלות קצב. בעומס שיא של 1,200 בקשות לדקה, כיסוי של 20 זוגות מסחר לא הספיק, והנתונים הגיעו באיחור של יותר משנייה. העברנו אותו ל-WebSocket, והשגנו זמן השהיה של 2 אלפיות שנייה ושלמות נתונים של 99.99%. מקרים כאלה הם הנורמה: מגבלות API, ניתוקי חיבור, חוסר עקביות בפורמטים. איסוף עסקאות בזמן אמת הוא אתגר הנדסי שאנחנו פותרים מקצה לקצה. אנחנו מתמודדים עם עד 100,000 עסקאות בשנייה על VPS יחיד. צרו קשר כדי להעריך את הפרויקט שלכם.
למה WebSocket הוא האפשרות היחידה לזמן אמת
סקירת REST מוסיפה זמן השהיה של 500–2000 אלפיות שנייה ומפספסת עסקאות בעומס שיא. WebSocket מספק סטרימינג עם זמן השהיה של 1–50 אלפיות שנייה. תיעוד WebSocket של Binance ממליץ על עד 300 סטרימים לכל חיבור. אנחנו משתמשים במספר חיבורים כדי לכסות את כל הזוגות.
ניהול מספר בורסות
CCXT Pro מספק ממשק watch_trades אחיד ל-30+ בורסות עם חיבור מחדש אוטומטי. הקוד שלהלן מתחבר לכל CEX בכמה שורות:
import ccxt.pro as ccxtpro import asyncio async def collect_trades(exchange_id: str, symbols: list[str], queue: asyncio.Queue): exchange = getattr(ccxtpro, exchange_id)({ 'enableRateLimit': True, 'options': {'tradesLimit': 1000}, }) try: while True: try: trades = await exchange.watch_trades_for_symbols(symbols) for trade in trades: await queue.put({ 'exchange': exchange_id, 'symbol': trade['symbol'], 'id': trade['id'], 'price': trade['price'], 'amount': trade['amount'], 'side': trade['side'], 'timestamp': trade['timestamp'], }) except Exception as e: print(f'Error {exchange_id}: {e}, reconnecting...') await asyncio.sleep(1) finally: await exchange.close() CCXT Pro מהיר פי 10 לפיתוח מאשר כתיבת מחברי WebSocket מותאמים לכל בורסה. לצורך קנה מידה, אנחנו משתמשים באשכולות: מספר מופעים המחלקים עומס בין קבוצות שונות של זוגות מסחר. זה מתמודד עם אלפי זוגות ללא אובדן ביצועים.
איך לאסוף עסקאות מ-DEX?
ב-DEX, עסקאות הן אירועי חוזה חכם. שתי שיטות עיקריות: תת-גרף The Graph — נתונים מוכנים דרך GraphQL. עבור Uniswap V3:
{ swaps( first: 100 orderBy: timestamp orderDirection: desc where: { pool: "0x8ad599c3a0ff1de082011efddc58f1908eb6e6d8" } ) { id timestamp amount0 amount1 sqrtPriceX96 tick transaction { id } } } זמן השהיה: 1–5 דקות מהכללת הבלוק. ניטור RPC ישיר — הרשמה לאירועי import ccxt.pro as ccxtpro import asyncio async def collect_trades(exchange_id: str, symbols: list[str], queue: asyncio.Queue): exchange = getattr(ccxtpro, exchange_id)({ 'enableRateLimit': True, 'options': {'tradesLimit': 1000}, }) try: while True: try: trades = await exchange.watch_trades_for_symbols(symbols) for trade in trades: await queue.put({ 'exchange': exchange_id, 'symbol': trade['symbol'], 'id': trade['id'], 'price': trade['price'], 'amount': trade['amount'], 'side': trade['side'], 'timestamp': trade['timestamp'], }) except Exception as e: print(f'Error {exchange_id}: {e}, reconnecting...') await asyncio.sleep(1) finally: await exchange.close() דרך { swaps( first: 100 orderBy: timestamp orderDirection: desc where: { pool: "0x8ad599c3a0ff1de082011efddc58f1908eb6e6d8" } ) { id timestamp amount0 amount1 sqrtPriceX96 tick transaction { id } } } . שלבים:
- התחברות לצומת RPC דרך WebSocket.
- הרשמה לאירוע
Swapעבור הבריכה. - פענוח
eth_subscribeלמחיר. - עיבוד ואחסון הנתונים.
דוגמה ב-TypeScript:
import { createPublicClient, webSocket, parseAbiItem } from 'viem'; const SWAP_EVENT = parseAbiItem( 'event Swap(address indexed sender, address indexed recipient, int256 amount0, int256 amount1, uint160 sqrtPriceX96, uint128 liquidity, int24 tick)' ); client.watchContractEvent({ address: UNISWAP_V3_POOL, event: SWAP_EVENT, onLogs: (logs) => { for (const log of logs) { const { amount0, amount1, sqrtPriceX96 } = log.args; const price = sqrtPriceX96ToPrice(sqrtPriceX96, token0Decimals, token1Decimals); processSwap({ price, amount0, amount1, txHash: log.transactionHash }); } } }); המחיר ב-Uniswap V3 מאוחסן כ-Swap (נקודה קבועה Q64.96). פענוח:
function sqrtPriceX96ToPrice(sqrtPriceX96: bigint, d0: number, d1: number): number { const price = Number(sqrtPriceX96 ** 2n * BigInt(10 ** d0)) / Number(BigInt(2 ** 192) * BigInt(10 ** d1)); return price; } השוואת שיטות איסוף נתוני DEX
| שיטה | זמן השהיה | תפוקה | מורכבות יישום |
|---|---|---|---|
| תת-גרף The Graph | 1–5 דקות | גבוהה | נמוכה |
| ניטור RPC ישיר | ~500 אלפיות שנייה | בינונית | בינונית |
| ניתוח לוגים דרך eth_getLogs | ~5 שניות | נמוכה | גבוהה |
איסוף עסקאות CEX לעומת DEX
| מאפיין | CEX | DEX |
|---|---|---|
| סוג נתונים | API מרכזי | אירועי שרשרת |
| זמן השהיה | 1–50 אלפיות שנייה (WebSocket) | 500 אלפיות שנייה – 5 דקות |
| אמינות | גבוהה (מגבלות IP) | תלוי בצומת |
| מורכבות אינטגרציה | בינונית (CCXT) | גבוהה (פענוח) |
איך לעקוף מגבלות קצב וחסימות?
Binance: 1,200 בקשות לדקה לכל IP עבור REST, WebSocket עד 300 סטרימים לכל חיבור. אנחנו משתמשים במספר חיבורים עם 300 זוגות כל אחד.
ל-Bybit ו-OKX יש מגבלות דומות. Bybit מנתקת WebSocket בחוסר פעילות—שלחו פינג כל 20 שניות. סיבוב IP עובד עבור REST אך לא עבור WebSocket. לתדירות גבוהה, אנחנו משתמשים במספר VPS במרכזי נתונים שונים.
אנחנו גם מגדירים דחיסה ומאגרי חיבורים כדי להפחית עומס. הניסיון של הצוות שלנו (50+ אינטגרציות) מבטיח יציבות גם בנפחי שיא.
תצורה לדוגמה לטיפול במגבלות קצב
# Настройка CCXT Pro с контролем лимитов exchange = ccxtpro.binance({ 'enableRateLimit': True, 'rateLimit': 1000, 'options': { 'tradesLimit': 1000, 'watchTrades': {'limit': 100}, }, }) לצורך אשכולות, אנחנו מריצים מספר מופעים כאלה, כל אחד עם קבוצת סמלים משלו.
איך לנרמל ולאחסן עסקאות?
אנחנו מאחדים את כל העסקאות לסכמה אחת המחולקת לפי יום:
CREATE TABLE trades ( id BIGSERIAL PRIMARY KEY, exchange VARCHAR(50) NOT NULL, symbol VARCHAR(30) NOT NULL, trade_id VARCHAR(100), price NUMERIC(30, 10) NOT NULL, quantity NUMERIC(30, 10) NOT NULL, side CHAR(4) NOT NULL, ts TIMESTAMPTZ NOT NULL, received_at TIMESTAMPTZ DEFAULT NOW() ) PARTITION BY RANGE (ts); CREATE INDEX ON trades (exchange, symbol, ts DESC); CREATE INDEX ON trades (symbol, ts DESC); TimescaleDB מפשטת זאת עם sqrtPriceX96 לאגרגציות OHLCV. שימוש ב-TimescaleDB מפחית עלויות אחסון בכ-30% בהשוואה ל-PostgreSQL בזכות דחיסה וחלוקה.
הסרת כפילויות
במהלך חיבור מחדש של WebSocket, השרת שולח מחדש את N העסקאות האחרונות. אילוץ ייחודי על import { createPublicClient, webSocket, parseAbiItem } from 'viem'; const SWAP_EVENT = parseAbiItem( 'event Swap(address indexed sender, address indexed recipient, int256 amount0, int256 amount1, uint160 sqrtPriceX96, uint128 liquidity, int24 tick)' ); client.watchContractEvent({ address: UNISWAP_V3_POOL, event: SWAP_EVENT, onLogs: (logs) => { for (const log of logs) { const { amount0, amount1, sqrtPriceX96 } = log.args; const price = sqrtPriceX96ToPrice(sqrtPriceX96, token0Decimals, token1Decimals); processSwap({ price, amount0, amount1, txHash: log.transactionHash }); } } }); מונע כפילויות.
מה כלול בעבודה שלנו
אנחנו מספקים:
- ארכיטקטורת מערכת לאיסוף עסקאות (בחירת סטרימים, אינטגרציית CEX + DEX)
- קוד Python/TypeScript עם טיפול בשגיאות וחיבור מחדש אוטומטי
- נירמול נתונים לסכמה אחידה
- פריסה על VPS/Kubernetes עם ניטור
- תיעוד והדרכת צוות
- 30 ימי תמיכה לאחר השחרור
עלות פיתוח מערכת איסוף עסקאות תלויה במספר הבורסות וזוגות המסחר. עבור סט טיפוסי של 3–5 בורסות ו-10–20 זוגות, היא נקבעת לאחר ניתוח מפורט. עם ניסיון של למעלה מ-10 שנים בפיתוח בלוקצ'יין, 50+ אינטגרציות בורסות ומהנדסים מוסמכים, אנחנו מבטיחים איסוף עסקאות אמין בכל עומס. צרו קשר—נכין ארכיטקטורה מותאמת לנפחים שלכם.







