גרידת נתוני ספר הזמנות בזמן אמת
תארו לעצמכם שבוט המסחר שלכם מבצע עסקה במחיר שכבר השתנה – ספר ההזמנות אינו מסונכרן עקב ניתוק WebSocket. הפסדים מאירוע יחיד כזה יכולים להגיע ל-2-3% מההון. נתקלנו בכך בפרויקטים מוקדמים ופיתחנו מערכת שמבטלת תקלות כאלה.
נתוני ספר ההזמנות קריטיים לשלושה תרחישים: בניית בוט מסחר, יצירת מצרף נזילות וניטור שוק. בכל אחד מהם, האתגר המשותף הוא השגת נתונים יציבה בתדירות גבוהה ללא אובדן ועם זמן השהיה מינימלי. בחירת פרוטוקול שגויה או חוסר טיפול בחיבור מחדש מובילים לחוסר סנכרון בספר ההזמנות ולעסקאות לא רווחיות. במהלך השנים, יישמנו למעלה מ-30 אינטגרציות עם בורסות שונות. אנו מבטיחים איסוף נתונים יציב וזמן השהיה נמוך.
למה WebSocket על פני REST?
סקירת REST (GET /api/v3/depth?symbol=BTCUSDT) היא הבחירה השגויה לספר הזמנות בזמן אמת. בשווקים פעילים, ספר ההזמנות מתעדכן 10–100 פעמים בשנייה. סקירה פעם בשנייה נותנת נתונים מיושנים ומכבידה על מגבלות קצב ה-API. הגישה הנכונה היא זרמי WebSocket עם עדכונים מצטברים.
| פרמטר | סקירת REST | זרם WebSocket |
|---|---|---|
| זמן השהיה | שנייה+ (מרווח סקירה) | 10-100 אלפיות שנייה (מונע אירועים) |
| עומס API | גבוה (בקשות כל שנייה) | נמוך (חיבור יחיד) |
| רעננות נתונים | מיושן מיד | תמיד במצב העדכני ביותר |
| קנה מידה | בעיות עם מכשירים מרובים | עד 1024 זרמים למפתח |
רוב הבורסות המרכזיות הגדולות (Binance, Bybit, OKX) פועלות לפי אותה תכנית:
- קבלת תמונת מצב דרך REST (ספר הזמנות מלא ברגע הנוכחי)
- הרשמה לזרם עדכוני WebSocket
- החלת עדכונים על תמונת המצב, תוך שמירה על עותק מקומי של ספר ההזמנות
import asyncio, json, aiohttp from sortedcontainers import SortedDict class OrderBook: def __init__(self): self.bids = SortedDict(lambda x: -x) # убывающий порядок self.asks = SortedDict() self.last_update_id = 0 def apply_update(self, bids: list, asks: list, update_id: int): if update_id <= self.last_update_id: return # устаревшее обновление, игнорируем for price, qty in bids: price, qty = float(price), float(qty) if qty == 0: self.bids.pop(price, None) # удалить уровень else: self.bids[price] = qty for price, qty in asks: price, qty = float(price), float(qty) if qty == 0: self.asks.pop(price, None) else: self.asks[price] = qty self.last_update_id = update_id @property def best_bid(self) -> tuple[float, float] | None: if self.bids: price = self.bids.keys()[0] return price, self.bids[price] return None @property def best_ask(self) -> tuple[float, float] | None: if self.asks: price = self.asks.keys()[0] return price, self.asks[price] return None איך לסנכרן מחדש את ספר ההזמנות לאחר ניתוק?
כאשר חיבור נופל או חבילות אובדות, הסיכון לחוסר עקביות גבוה. אנו משתמשים בטכניקות הבאות:
- אגירת עדכונים עד לקבלת תמונת מצב (כפי שמוצג לעיל)
- בדיקת מזהה הרצף של כל עדכון: אם
import asyncio, json, aiohttp from sortedcontainers import SortedDict class OrderBook: def __init__(self): self.bids = SortedDict(lambda x: -x) # убывающий порядок self.asks = SortedDict() self.last_update_id = 0 def apply_update(self, bids: list, asks: list, update_id: int): if update_id <= self.last_update_id: return # устаревшее обновление, игнорируем for price, qty in bids: price, qty = float(price), float(qty) if qty == 0: self.bids.pop(price, None) # удалить уровень else: self.bids[price] = qty for price, qty in asks: price, qty = float(price), float(qty) if qty == 0: self.asks.pop(price, None) else: self.asks[price] = qty self.last_update_id = update_id @property def best_bid(self) -> tuple[float, float] | None: if self.bids: price = self.bids.keys()[0] return price, self.bids[price] return None @property def best_ask(self) -> tuple[float, float] | None: if self.asks: price = self.asks.keys()[0] return price, self.asks[price] return Noneאינו תואם למצופה, זרוק את החבילה ובקש תמונת מצב חדשה - השהיה אקספוננציאלית בחיבור מחדש עם תקרה של 60 שניות
- ניטור זמן השהיה והתראה כאשר חריגה מהסף (לדוגמה, >500 אלפיות שנייה)
async def connect_ws_with_retry(url: str, handler, max_retries=10): for attempt in range(max_retries): try: async with websockets.connect(url, ping_interval=20) as ws: async for message in ws: await handler(message) except (websockets.exceptions.ConnectionClosed, Exception) as e: wait = min(2 ** attempt, 60) # max 60 секунд logging.warning(f"WS disconnected: {e}, retry in {wait}s") await asyncio.sleep(wait) פרטי עבודה עם זרם העומק של Binance
Binance היא הבקשה הנפוצה ביותר. יש להם שתי גרסאות זרם:
-
update_id— עדכונים כל 100 אלפיות שנייה או 1000 אלפיות שנייה (פרמטרasync def connect_ws_with_retry(url: str, handler, max_retries=10): for attempt in range(max_retries): try: async with websockets.connect(url, ping_interval=20) as ws: async for message in ws: await handler(message) except (websockets.exceptions.ConnectionClosed, Exception) as e: wait = min(2 ** attempt, 60) # max 60 секунд logging.warning(f"WS disconnected: {e}, retry in {wait}s") await asyncio.sleep(wait)) -
btcusdt@depth— 20 רמות עליונות כל 100 אלפיות שנייה (ללא עדכונים מצטברים, תמיד מלא)
לספר הזמנות מלא עם תיקונים:
async def maintain_binance_orderbook(symbol: str): ob = OrderBook() buffer = [] # буфер обновлений до получения snapshot async def handle_ws_message(msg): data = json.loads(msg) # Накапливаем обновления ПОКА не получим snapshot if ob.last_update_id == 0: buffer.append(data) return # Binance: обновление валидно если U <= lastUpdateId+1 <= u if data['U'] <= ob.last_update_id + 1 <= data['u']: ob.apply_update(data['b'], data['a'], data['u']) # Запускаем WS ws_task = asyncio.create_task(connect_ws( f"wss://stream.binance.com:9443/ws/{symbol.lower()}@depth@100ms", handle_ws_message )) # Получаем snapshot (немного ждём чтобы буфер накопился) await asyncio.sleep(0.5) async with aiohttp.ClientSession() as session: async with session.get( f"https://api.binance.com/api/v3/depth", params={"symbol": symbol.upper(), "limit": 1000} ) as resp: snapshot = await resp.json() # Инициализируем стакан из snapshot for price, qty in snapshot['bids']: ob.bids[float(price)] = float(qty) for price, qty in snapshot['asks']: ob.asks[float(price)] = float(qty) ob.last_update_id = snapshot['lastUpdateId'] # Применяем буферизованные обновления for update in buffer: if update['u'] > ob.last_update_id: ob.apply_update(update['b'], update['a'], update['u']) await ws_task נקודה קריטית: אם עדכון הוחמץ (פער ברצף @depth@100ms → btcusdt@depth20) ספר ההזמנות הופך לבלתי מסונכרן. יש צורך בלוגיקת סנכרון מחדש: זיהוי הפער ואתחול מחדש מתמונת מצב חדשה.
מקרה בוחן: צירוף Binance ו-Bybit לארביטראז'
לארביטראז' בין-בורסתי, יש צורך לתחזק ספרי הזמנות של מספר בורסות במקביל. הנה דוגמה למצרף שמוצא את המחיר הטוב ביותר:
EXCHANGES = { "binance": BinanceOrderBook, "bybit": BybitOrderBook, "okx": OKXOrderBook, } async def run_aggregator(symbol: str): books = {name: cls(symbol) for name, cls in EXCHANGES.items()} tasks = [book.run() for book in books.values()] await asyncio.gather(*tasks) def get_best_price_across_exchanges(books: dict[str, OrderBook]) -> dict: best_bids = [(name, *ob.best_bid) for name, ob in books.items() if ob.best_bid] best_asks = [(name, *ob.best_ask) for name, ob in books.items() if ob.best_ask] best_bids.sort(key=lambda x: x[1], reverse=True) best_asks.sort(key=lambda x: x[1]) return { "best_bid": {"exchange": best_bids[0][0], "price": best_bids[0][1], "qty": best_bids[0][2]}, "best_ask": {"exchange": best_asks[0][0], "price": best_asks[0][1], "qty": best_asks[0][2]}, "spread": best_asks[0][1] - best_bids[0][1] } צירוף בזמן אמת מאפשר לראות את ההצעה והביקוש הטובים ביותר בכל הבורסות. זהו הבסיס לאסטרטגיות ארביטראז' ולבניית ספר הזמנות מאוחד.
אחסון נתונים: TimescaleDB לעומת מבוסס קבצים
לבדיקות חוזרות או ביקורת, עדיף לאחסן את זרם העדכונים ולא רק תמונות מצב. עדכוני ספר הזמנות L2 מייצרים נפח גדול: עבור BTC/USDT ב-Binance כ-100MB לשעה של נתונים לא דחוסים.
| קריטריון | TimescaleDB | מבוסס קבצים (Parquet) |
|---|---|---|
| שאילתות בזמן אמת | כן (SQL) | לא (אנליטיקה בלבד) |
| דחיסה | אוטומטית | ניתנת להגדרה (lz4) |
| השמעה חוזרת של זרם | דורש עיבוד נוסף | קריאה ישירה |
| אינטגרציית Kafka | כן | לא זמין |
אנו ממליצים להשתמש ב-TimescaleDB לאחסון לטווח ארוך וב-Parquet לאנליטיקה. במידת הצורך, אנו משלבים את הזרם לתוך Kafka עבור מערכות במורד הזרם.
# Запись в бинарный формат через msgpack import msgpack, lz4.frame def serialize_update(update: dict) -> bytes: packed = msgpack.packb(update, use_bin_type=True) return lz4.frame.compress(packed) # TimescaleDB для time-series хранения # Гипертаблица автоматически партиционирует по времени CREATE TABLE ob_updates ( time TIMESTAMPTZ NOT NULL, exchange TEXT NOT NULL, symbol TEXT NOT NULL, side CHAR(1) NOT NULL, -- 'b' или 'a' price NUMERIC NOT NULL, quantity NUMERIC NOT NULL ); SELECT create_hypertable('ob_updates', 'time'); מה כלול בעבודה?
בהזמנת מערכת גרידת ספר הזמנות, אתם מקבלים:
- פתרון ארכיטקטוני עם בחירת פרוטוקול ואסטרטגיית סנכרון מחדש
- קוד מקור ב-Python עם asyncio ותיעוד פריסה
- לוחות מחוונים של Grafana לניטור זמן השהיה ושגיאות
- הבטחת יציבות של 99.9% ושבועיים של תמיכה לאחר הפריסה
לוח זמנים לפיתוח: 2 עד 4 שבועות תלוי במספר הבורסות והארכיטקטורה הנדרשת. העלות מחושבת באופן אישי. צרו קשר כדי לדון בפרויקט שלכם – נמצא פתרון אופטימלי.
טעויות נפוצות בגרידת ספר הזמנות
- התעלמות ממזהה הרצף וחוסר סנכרון מחדש לאחר פער
- שימוש בסקירת REST במקום WebSocket (מוביל לעיכובים ומגבלות קצב)
- סדר שגוי: תמונת מצב תחילה, ולאחר מכן הרשמה לעדכונים
- אין אגירת עדכונים לפני תמונת המצב (החבילות הראשונות אובדות)
- טיפול שגוי בחיבור מחדש ללא השהיה אקספוננציאלית
הזמינו פיתוח של מערכת גרידת ספר הזמנות עם הבטחת יציבות וזמן השהיה נמוך. הניסיון שלנו במסחר בתדירות גבוהה מאפשר לנו ליצור פתרונות שאינם מאבדים נתונים ואינם יוצאים מסנכרון אפילו בעומסי שיא.







