ספר הזמנות מלא מכיל את כל מידע הנזילות בבורסה. איסוף שלו, נירמולו והפיכתו לתכונות עבור למידת מכונה הוא אתגר הנדסי לא טריוויאלי. בנינו צינור עיבוד ברמת ייצור עבור Binance, Bybit ו-OKX שמעבד עד 10,000 עדכונים בשנייה. הניסיון שלנו כולל אינטגרציה עם יותר מ-15 בורסות קריפטו ואחסון של כ-5 TB של נתונים בחודש. ספר הזמנות L2 מלא מתאר כל רמת מחיר עם נפח — זה הבסיס לבניית תחזיות לטווח קצר. איסוף יציב תחת עומסי שיא ועקביות של תמונות מצב מובטחים.
לקוחות מגיעים לעיתים קרובות עם זרמי WebSocket גולמיים, לא בטוחים כיצד לסנכרן את זרם ההפרשים עם תמונת מצב REST. שגיאת off-by-one גורמת לספר לסטות, מה שמוביל לאותות שגויים. אנו פותרים זאת ברמת ארכיטקטורת האספן.
בעיות שאנו פותרים
- נפח נתונים. ספר הזמנות L2 מלא ב-Binance מכיל 5000 רמות בכל צד. עם עדכונים כל 100 אלפיות שנייה, זה מייצר עשרות גיגה-בייט ביום. אחסון נאיבי ב-PostgreSQL יהרוג את הביצועים.
- תנאי מרוץ. זרם ההפרשים של WebSocket מגיע באופן אסינכרוני. ללא סנכרון עם תמונת המצב של REST, הספר סוטה — מחירים הולכים לרמות לא קיימות.
- פורמט נתונים. כל בורסה מספקת את ספר ההזמנות אחרת: Binance משתמשת במערכים מקוננים, Coinbase משתמשת ב-JSON עם מפתחות שונים. יש צורך בממשק אחיד.
כיצד לסנכרן זרם הפרשים של WebSocket עם תמונת מצב REST?
האלגוריתם פשוט: פתחו WebSocket, קבלו את זרם ההפרשים הראשון, בקשו מיד תמונת מצב REST מלאה. לאחר מכן החילו כל עדכון על ספר ההזמנות המקומי. השתמשו ב-lastUpdateId לשליטה: החילו רק הודעות עם u > lastUpdateId. אם הרצף נשבר — בקשו שוב תמונת מצב. גישה זו מבטלת סטיית ספר גם תחת תנודתיות גבוהה.
כיצד לאסוף ספר הזמנות דרך WebSocket: אלגוריתם שלב אחר שלב
- יצירת חיבור: דרך
wss://stream.binance.com:9443/ws/btcusdt@depth@100ms(אנלוגי לבורסות אחרות). - תמונת מצב REST ראשונית: סנכרון דרך
updateIdכדי להבטיח עקביות. - עדכונים מצטברים: כל הודעת זרם הפרשים מוחלת על מצב הספר הנוכחי.
- שמירת תמונות מצב: בתדירות נתונה (כל עדכון N), קבעו את המצב המלא להנדסת תכונות עתידית.
דוגמת קוד אספן:
import asyncio
import websockets
import json
from collections import deque
class OrderBookCollector:
def __init__(self, symbol, max_depth=100):
self.symbol = symbol
self.bids = {}
self.asks = {}
self.max_depth = max_depth
self.snapshots = deque(maxlen=10000)
async def connect_binance(self):
url = f"wss://stream.binance.com:9443/ws/{self.symbol.lower()}@depth@100ms"
async with websockets.connect(url) as ws:
await self.fetch_snapshot()
async for msg in ws:
data = json.loads(msg)
self.process_diff_update(data)
if len(self.snapshots) % 10 == 0:
self.save_snapshot()
def process_diff_update(self, data):
for bid_level in data.get('b', []):
price, qty = float(bid_level[0]), float(bid_level[1])
if qty == 0:
self.bids.pop(price, None)
else:
self.bids[price] = qty
for ask_level in data.get('a', []):
price, qty = float(ask_level[0]), float(ask_level[1])
if qty == 0:
self.asks.pop(price, None)
else:
self.asks[price] = qty
def get_features(self, n_levels=20):
sorted_bids = sorted(self.bids.items(), reverse=True)[:n_levels]
sorted_asks = sorted(self.asks.items())[:n_levels]
if not sorted_bids or not sorted_asks:
return None
mid_price = (sorted_bids[0][0] + sorted_asks[0][0]) / 2
features = {}
for i, (price, qty) in enumerate(sorted_bids[:10]):
features[f'bid_qty_{i}'] = qty
features[f'bid_dist_{i}'] = (mid_price - price) / mid_price
for i, (price, qty) in enumerate(sorted_asks[:10]):
features[f'ask_qty_{i}'] = qty
features[f'ask_dist_{i}'] = (price - mid_price) / mid_price
bid_vol_n = sum(qty for _, qty in sorted_bids[:5])
ask_vol_n = sum(qty for _, qty in sorted_asks[:5])
features['obi_5'] = (bid_vol_n - ask_vol_n) / (bid_vol_n + ask_vol_n + 1e-8)
bid_vol_20 = sum(qty for _, qty in sorted_bids[:20])
ask_vol_20 = sum(qty for _, qty in sorted_asks[:20])
features['obi_20'] = (bid_vol_20 - ask_vol_20) / (bid_vol_20 + ask_vol_20 + 1e-8)
features['wmid'] = (sorted_bids[0][0] * sorted_asks[0][1] + sorted_asks[0][0] * sorted_bids[0][1]) / (sorted_bids[0][1] + sorted_asks[0][1])
features['spread'] = (sorted_asks[0][0] - sorted_bids[0][0]) / mid_price
for n in [5, 10, 20]:
bid_depth = sum(qty for _, qty in sorted_bids[:n])
ask_depth = sum(qty for _, qty in sorted_asks[:n])
features[f'depth_ratio_{n}'] = bid_depth / max(ask_depth, 1e-8)
return features
מדוע ClickHouse הוא אחסון אופטימלי לספר הזמנות?
ספר הזמנות L2 מלא הוא עצום. ClickHouse מהיר פי 10 מ-PostgreSQL באגרגציות עמודות. לפי תיעוד ClickHouse, DBMS עמודי מספק דחיסה של עד פי 10 ומהירויות כתיבה של למעלה ממיליון שורות בשנייה. השוו:
| DBMS | מהירות כתיבה (שורות/שנייה) | דחיסה | אגרגציות מבוססות זמן |
|---|---|---|---|
| PostgreSQL | ~100,000 | 2-5x | איטי |
| TimescaleDB | ~200,000 | 3-6x | בינוני |
| ClickHouse | ~1,000,000 | 5-10x | מהיר |
דוגמת סכמה עם TTL אוטומטי:
CREATE TABLE order_book_snapshots (
timestamp DateTime64(3),
symbol LowCardinality(String),
exchange LowCardinality(String),
bid_price_0 Float32,
bid_qty_0 Float32,
bid_price_1 Float32,
bid_qty_1 Float32,
-- ... до bid_price_19, bid_qty_19
ask_price_0 Float32,
ask_qty_0 Float32,
-- ...
spread Float32,
obi_5 Float32,
obi_20 Float32
) ENGINE = MergeTree()
PARTITION BY toYYYYMMDD(timestamp)
ORDER BY (symbol, timestamp)
TTL timestamp + INTERVAL 90 DAY;חיסכון בתשתית בעת שימוש ב-ClickHouse מגיע ל-70% בזכות דחיסה — זה כ-$20,000 בשנה עבור פרויקט עם 5 TB של נתונים. עבור פרויקטים גדולים, החיסכון יכול להגיע עד $30,000 בשנה.
הנדסת תכונות מספר הזמנות
בהתבסס על תמונות מצב שנאספו, אנו בונים תכונות. בסיסיות: OBI (חוסר איזון בספר הזמנות), מרווח, עומק. נוספות: ממוצעים נעים של OBI, התנודתיות שלו, זרימת הזמנות מצטברת (COF).
def engineer_orderbook_features(snapshots_df, window_sizes=[10, 50, 100]):
features = snapshots_df.copy()
for window in window_sizes:
features[f'obi_5_ma_{window}'] = features['obi_5'].rolling(window).mean()
features[f'obi_5_delta_{window}'] = features['obi_5'].diff(window)
features[f'obi_5_std_{window}'] = features['obi_5'].rolling(window).std()
features['cof'] = features['obi_5'].cumsum()
features['cof_ma'] = features['cof'].rolling(100).mean()
features['cof_deviation'] = features['cof'] - features['cof_ma']
features['spread_ma'] = features['spread'].rolling(50).mean()
features['spread_ratio'] = features['spread'] / features['spread_ma']
features['depth_change'] = features['depth_ratio_10'].diff(10)
return features
כיצד להעריך איכות חיזוי מחיר אמצע?
עבור חיזוי מחיר אמצע לטווח קצר (לאחר N עדכוני ספר) אנו משתמשים בדיוק, דיוק חיובי ו-F1-score לסיווג בינארי של כיוון. קוד להכנת נתוני אימון:
def create_training_data(snapshots_df, prediction_horizon=10):
features = engineer_orderbook_features(snapshots_df)
future_mid = snapshots_df['mid_price'].shift(-prediction_horizon)
current_mid = snapshots_df['mid_price']
target = np.sign(future_mid - current_mid)
valid_mask = features.notna().all(axis=1) & target.notna()
return features[valid_mask], target[valid_mask] טעויות נפוצות בפיתוח צינור ספר הזמנות
אפילו צוותים מנוסים עושים טעויות: התעלמות מהטיית ספר במהלך תנודתיות גבוהה, טיפול שגוי באירועי import asyncio import websockets import json from collections import deque class OrderBookCollector: def __init__(self, symbol, max_depth=100): self.symbol = symbol self.bids = {} self.asks = {} self.max_depth = max_depth self.snapshots = deque(maxlen=10000) async def connect_binance(self): url = f"wss://stream.binance.com:9443/ws/{self.symbol.lower()}@depth@100ms" async with websockets.connect(url) as ws: await self.fetch_snapshot() async for msg in ws: data = json.loads(msg) self.process_diff_update(data) if len(self.snapshots) % 10 == 0: self.save_snapshot() def process_diff_update(self, data): for bid_level in data.get('b', []): price, qty = float(bid_level[0]), float(bid_level[1]) if qty == 0: self.bids.pop(price, None) else: self.bids[price] = qty for ask_level in data.get('a', []): price, qty = float(ask_level[0]), float(ask_level[1]) if qty == 0: self.asks.pop(price, None) else: self.asks[price] = qty def get_features(self, n_levels=20): sorted_bids = sorted(self.bids.items(), reverse=True)[:n_levels] sorted_asks = sorted(self.asks.items())[:n_levels] if not sorted_bids or not sorted_asks: return None mid_price = (sorted_bids[0][0] + sorted_asks[0][0]) / 2 features = {} for i, (price, qty) in enumerate(sorted_bids[:10]): features[f'bid_qty_{i}'] = qty features[f'bid_dist_{i}'] = (mid_price - price) / mid_price for i, (price, qty) in enumerate(sorted_asks[:10]): features[f'ask_qty_{i}'] = qty features[f'ask_dist_{i}'] = (price - mid_price) / mid_price bid_vol_n = sum(qty for _, qty in sorted_bids[:5]) ask_vol_n = sum(qty for _, qty in sorted_asks[:5]) features['obi_5'] = (bid_vol_n - ask_vol_n) / (bid_vol_n + ask_vol_n + 1e-8) bid_vol_20 = sum(qty for _, qty in sorted_bids[:20]) ask_vol_20 = sum(qty for _, qty in sorted_asks[:20]) features['obi_20'] = (bid_vol_20 - ask_vol_20) / (bid_vol_20 + ask_vol_20 + 1e-8) features['wmid'] = (sorted_bids[0][0] * sorted_asks[0][1] + sorted_asks[0][0] * sorted_bids[0][1]) / (sorted_bids[0][1] + sorted_asks[0][1]) features['spread'] = (sorted_asks[0][0] - sorted_bids[0][0]) / mid_price for n in [5, 10, 20]: bid_depth = sum(qty for _, qty in sorted_bids[:n]) ask_depth = sum(qty for _, qty in sorted_asks[:n]) features[f'depth_ratio_{n}'] = bid_depth / max(ask_depth, 1e-8) return features , חוסר בדיקות עקביות לאחר חיבור מחדש. נתקלנו בפרויקט שבו בגלל הפרשים שהוחמצו הספר סטה ב-20% — המודל נתן אותות שגויים. הפתרון הוא הטמעת בדיקות checksum ושחזור תמונת מצב מלא אוטומטי בעת זיהוי חוסר עקביות.
מה כלול בפיתוח הצינור
- קוד מקור לאספן ולצינור (Python אסינכרוני).
- דמפי נתוני בדיקה לבדיקות אופליין.
- README עם דוגמאות שימוש מפורטות.
- מיגרציות סכמה של ClickHouse עם TTL.
- הכשרת הצוות שלך לשימוש בצינור.
שלבים ולוחות זמנים
| שלב | משך | תוצאה |
|---|---|---|
| אנליטיקה | 2-3 ימים | מפרט API, הערכות נפח |
| עיצוב | 2-3 ימים | סכמת אחסון, בחירת תכונות |
| יישום | 5-10 ימים | אספן, צינור, קוד |
| בדיקות | 3-5 ימים | סימולציה של 24 שעות, דוחות |
| פריסה | 2-3 ימים | Docker, ניטור |
צינור בסיסי לבורסה אחת עם מודל LightGBM לוקח בין 14 ל-30 ימי עבודה. עלויות הפיתוח מתחילות ב-$15,000. אנו מספקים הערכה מדויקת לאחר ביקורת חינם של הנתונים שלך. בקשו ניתוח ואנו נבחר את הארכיטקטורה האופטימלית לנפח ספר ההזמנות שלך. עם ניסיון של למעלה מ-5 שנים ויותר מ-50 פרויקטים שהושלמו, אנו מבטיחים אספקה אמינה. צרו קשר לייעוץ.







