אנו מפתחים צינורות עיבוד נתוני טיקים — תיעוד כל עסקה עם מחיר, נפח וצד. נרות OHLCV סטנדרטיים מאבדים את מבנה המיקרו של השוק: חוסר איזון בנזילות, עסקאות גדולות, זרימת קנייה/מכירה. ללא צינור איכותי, מודל ML מתאמן על רעש. לדוגמה, בפרויקט אחד עבור Binance (מהניסיון שלנו), העומס הגיע ל-300,000 טיקים בשנייה — ClickHouse טיפל בזה, בעוד PostgreSQL קרס ב-10,000. הניסיון שלנו של למעלה מחמש שנים מבטיח אמינות. צרו קשר — אנו מוכנים לתכנן וליישם צינור למשימות שלכם.
למה נתוני טיקים חשובים יותר מ-OHLCV עבור ML
בעת צבירה לנרות של דקה אחת, עד 80% מהמידע אובד: אינכם רואים כיצד העסקאות מתחלקות בתוך המרווח, האם היה קפיצת נפח, או מי היה התוקפן. ברי נפח, ברי דולר וברי חוסר איזון משמרים את האותות הללו. מודלי ML המאומנים על טיקים מראים דיוק גבוה ב-15–20% במשימות חיזוי כיוון מחיר.
בעיות שאנו פותרים
- עומסים גבוהים. בורסות מייצרות עד 500,000 טיקים בשנייה. מסדי נתונים סטנדרטיים אינם יכולים להתמודד עם קצבי הכנסה כאלה.
- השהיה. עבור אסטרטגיות HFT, העיכוב מקבלת הטיק לאות לא יעלה על 10 אלפיות השנייה.
- אחסון. נתוני טיקים לשנה מסתכמים בעשרות טרה-בייט. יש צורך בחלוקה, TTL ודחיסה יעילה. החיסכון בתשתית ClickHouse יכול להגיע ל-50% בהשוואה למסדי נתונים יחסיים מסורתיים.
- מגוון ברים. ברי זמן אינם אחידים בתקופות פעילות נמוכה. ברי נפח/דולר/חוסר איזון מתאימים את עצמם לפעילות השוק.
איך אנחנו עושים את זה: טכנולוגיה ומקרה בוחן
בפרויקט אחד עבור Binance (מהניסיון שלנו), בנינו צינור שאוסף עסקאות מצטברות דרך WebSocket, מאגר אותן בזיכרון, ומכניס באופן אסינכרוני ל-ClickHouse.
import asyncio
import websockets
import json
from datetime import datetime
import asyncpg
class TickDataCollector:
def __init__(self, symbol, db_pool):
self.symbol = symbol
self.db_pool = db_pool
self.buffer = []
self.buffer_size = 1000
async def connect_binance_trades(self):
url = f"wss://stream.binance.com:9443/ws/{self.symbol.lower()}@aggTrade"
async with websockets.connect(url, ping_interval=20) as ws:
async for msg in ws:
trade = json.loads(msg)
tick = {
'symbol': self.symbol,
'timestamp': datetime.fromtimestamp(trade['T'] / 1000),
'price': float(trade['p']),
'quantity': float(trade['q']),
'is_buyer_maker': trade['m'],
'trade_id': trade['a']
}
self.buffer.append(tick)
if len(self.buffer) >= self.buffer_size:
await self.flush_to_db()
async def flush_to_db(self):
async with self.db_pool.acquire() as conn:
await conn.executemany(
"""INSERT INTO trades (symbol, timestamp, price, quantity, is_buyer_maker, trade_id) VALUES ($1, $2, $3, $4, $5, $6)""",
[(t['symbol'], t['timestamp'], t['price'], t['quantity'], t['is_buyer_maker'], t['trade_id']) for t in self.buffer]
)
self.buffer.clear()
"}האחסון מאורגן ב-ClickHouse עם מנוע MergeTree, חלוקה יומית ו-TTL של 365 ימים. זה מספק דחיסה יעילה (פי 10 בהשוואה ל-CSV) וקצב הכנסה גבוה.
CREATE TABLE trades (
timestamp DateTime64(3),
symbol LowCardinality(String),
price Float64,
quantity Float32,
is_buyer_maker UInt8,
trade_id UInt64
) ENGINE = MergeTree()
PARTITION BY toYYYYMMDD(timestamp)
ORDER BY (symbol, timestamp)
TTL timestamp + INTERVAL 365 DAY
SETTINGS index_granularity = 8192;ClickHouse מכניס 500K+ שורות בשנייה — פי 50 מהר יותר מ-PostgreSQL לעומסים כאלה. צבירות חודשיות מסתיימות תוך שניות. אנו מבטיחים שהצינור שלכם יתמודד עם כל פעילות שוק.
איך לבנות ברי נפח מטיקים: שלב אחר שלב
- התחברו ל-WebSocket של הבורסה כדי לקבל עסקאות מצטברות.
- צברו טיקים במאגר (לדוגמה, 1000 רשומות).
- כאשר הנפח שצוין מושג, סגרו את הבר ושמרו אותו ל-ClickHouse.
- השתמשו בפונקציית
import asyncio import websockets import json from datetime import datetime import asyncpg class TickDataCollector: def __init__(self, symbol, db_pool): self.symbol = symbol self.db_pool = db_pool self.buffer = [] self.buffer_size = 1000 async def connect_binance_trades(self): url = f"wss://stream.binance.com:9443/ws/{self.symbol.lower()}@aggTrade" async with websockets.connect(url, ping_interval=20) as ws: async for msg in ws: trade = json.loads(msg) tick = { 'symbol': self.symbol, 'timestamp': datetime.fromtimestamp(trade['T'] / 1000), 'price': float(trade['p']), 'quantity': float(trade['q']), 'is_buyer_maker': trade['m'], 'trade_id': trade['a'] } self.buffer.append(tick) if len(self.buffer) >= self.buffer_size: await self.flush_to_db() async def flush_to_db(self): async with self.db_pool.acquire() as conn: await conn.executemany( """INSERT INTO trades (symbol, timestamp, price, quantity, is_buyer_maker, trade_id) VALUES ($1, $2, $3, $4, $5, $6)""", [(t['symbol'], t['timestamp'], t['price'], t['quantity'], t['is_buyer_maker'], t['trade_id']) for t in self.buffer] ) self.buffer.clear()מהדוגמה למטה.
ברי נפח נסגרים כאשר נפח נתון מצטבר, לא במרווח זמן קבוע. זה מניב מספר אחיד של תצפיות ללא קשר לפעילות השוק.
def create_volume_bars(ticks_df, bar_volume=10):
"""Каждый бар = bar_volume единиц актива"""
bars = []
current_bar = {'open': None, 'high': -np.inf, 'low': np.inf, 'close': None, 'volume': 0, 'start_time': None}
for _, tick in ticks_df.iterrows():
if current_bar['open'] is None:
current_bar['open'] = tick['price']
current_bar['start_time'] = tick['timestamp']
current_bar['high'] = max(current_bar['high'], tick['price'])
current_bar['low'] = min(current_bar['low'], tick['price'])
current_bar['close'] = tick['price']
current_bar['volume'] += tick['quantity']
if current_bar['volume'] >= bar_volume:
bars.append(current_bar.copy())
current_bar = {'open': None, 'high': -np.inf, 'low': np.inf, 'close': None, 'volume': 0, 'start_time': None}
return pd.DataFrame(bars)באופן דומה, נבנים ברי דולר (לפי נפח USD) וברי חוסר איזון (לפי חוסר איזון קנייה/מכירה).
| סוג בר | קריטריון סגירה | מתי להשתמש |
|---|---|---|
| זמן | מרווח זמן | נזילות גבוהה, פעילות אחידה |
| נפח | נפח מצטבר | התאמה לקפיצות תנודתיות |
| דולר | נפח USD מצטבר | בלתי תלוי במחיר הנכס |
| חוסר איזון | חוסר איזון קנייה/מכירה | מציאת נקודות היפוך |
אילו יתרונות מספק הנדסת תכונות מטיקים?
תכונות המופקות מטיקים משפרות את איכות מודל ה-ML: חוסר איזון בזרימה, תדירות עסקאות, סטייה מ-VWAP, יחס עסקאות גדולות. בצינור ML בזמן אמת, תכונות אלה מחושבות על חלונות נעים.
def create_tick_features(ticks_df, window_ticks=[50, 200, 1000]):
features = []
for i in range(max(window_ticks), len(ticks_df)):
row_features = {}
for window in window_ticks:
window_data = ticks_df.iloc[i-window:i]
buy_vol = window_data[~window_data['is_buyer_maker']]['quantity'].sum()
sell_vol = window_data[window_data['is_buyer_maker']]['quantity'].sum()
row_features[f'flow_imbalance_{window}'] = (
(buy_vol - sell_vol) / (buy_vol + sell_vol + 1e-8)
)
row_features[f'trade_frequency_{window}'] = (
window / (window_data['timestamp'].max() - window_data['timestamp'].min()).total_seconds() + 1e-8
)
row_features[f'avg_trade_size_{window}'] = window_data['quantity'].mean()
row_features[f'large_trade_ratio_{window}'] = (
(window_data['quantity'] > window_data['quantity'].quantile(0.9)).mean()
)
vwap = (window_data['price'] * window_data['quantity']).sum() / window_data['quantity'].sum()
row_features[f'vwap_deviation_{window}'] = (
ticks_df.iloc[i]['price'] - vwap
) / vwap
features.append(row_features)
return pd.DataFrame(features)עסקאות גדולות (מעל האחוזון ה-99) מעידות לעיתים קרובות על פעילות מוסדית. ניתוח הכיוון שלהן מספק אות נוסף.
איך להבטיח השהיה <10 אלפיות השנייה?
ארכיטקטורת סטרימינג אמיתית:
Binance WebSocket → asyncio consumer → buffer → ClickHouse batch insert → Redis sorted set (last 10k ticks) → Feature calculator (sliding window) → ML inference → Signal output ההשהיה מטיק לאות היא מתחת ל-10 אלפיות השנייה. מושגת באמצעות קלט/פלט אסינכרוני, אגירת Redis ותכונות מחושבות מראש על חלונות זמן. החיסכון באשכול ClickHouse בהשוואה למסדי נתונים מסורתיים יכול להגיע ל-50%.
לפי תיעוד ClickHouse, קצב ההכנסה מגיע ל-500,000 שורות בשנייה תיעוד ClickHouse.
תהליך העבודה
| שלב | משך | תוצאה |
|---|---|---|
| אנליטיקה | 1–2 ימים | מסמך דרישות וסכימת נתונים |
| עיצוב | 2–3 ימים | בחירת טכנולוגיה, עיצוב סכימת DB, הגדרת סוג בר |
| יישום | 1–2 שבועות | קולט, מצברים, הנדסת תכונות, אינטגרציית צינור ML |
| בדיקות | 3–5 ימים | ולידציה על נתונים היסטוריים, מבחן עומס מהירות |
| פריסה | 2–3 ימים | פריסה באשכול שלכם (Docker/K8s), ניטור |
ציר זמן ותוצרים
גרסה בסיסית (סימבול יחיד, ClickHouse, Redis) — החל משבועיים. צינור מלא עם ברי נפח/דולר/חוסר איזון, הנדסת תכונות והסקה בזמן אמת — החל מ-4 שבועות. עלות הפרויקט משתנה; תמחור מדויק נקבע לאחר ניתוח. השקעה בצינור איכותי משתלמת באמצעות שיפור דיוק מודל ה-ML למסחר.
מה כלול:
- תיעוד ארכיטקטורה.
- קוד מקור של הצינור עם הערות.
- הגדרת ClickHouse, Redis ותור.
- אינטגרציה עם תשתית ה-ML שלכם.
- הדרכת צוות (2–3 שיחות).
- חודשיים של תמיכה לאחר הפריסה.
רשימת בדיקה לאימות הצינור
- בדקו את קצב ההכנסה: ClickHouse חייב להכניס לפחות 100K שורות בשנייה על ליבה אחת.
- ודאו ש-TTL מוגדר — בלעדיו, הדיסק מתמלא תוך חודש.
- הגדירו ניטור השהיה לכל שלב.
- בדקו התחברות אוטומטית מחדש של WebSocket בעת ניתוק.
- אמתו צבירות על נתונים היסטוריים — השוו עם ברי ייחוס.
טעויות נפוצות
- חלוקה עדינה מדי (לפי שעה): מספר גדול של מחיצות פוגע בביצועי ClickHouse. האופטימלי הוא לפי יום.
- התעלמות מ-TTL: ללא ניקוי נתונים אוטומטי, הדיסק מתמלא תוך חודש.
- שימוש בברי זמן לנכסים עם נזילות נמוכה: רוב הנרות יהיו ריקים.
אנו מפתחים צינורות נתוני טיקים במשך למעלה מחמש שנים, תוך יישום 30+ פרויקטים למסחר בקריפטו. קבלו ניתוח חינם של הנתונים שלכם והמלצות לאופטימיזציה של הצינור. צרו קשר — אנו נעריך את הפרויקט שלכם ונציע את הפתרון הטוב ביותר.







