Розробка pipeline обробки tick-даних для ML

Ми розробляємо пайплайни обробки tick-даних — записи кожної угоди з ціною, об'ємом та стороною. Стандартні OHLCV свічки втрачають мікроструктуру ринку: дисбаланс ліквідності, великі угоди, потік buy/sell. Без якісного пайплайну ML-модель навчається на шумі. Наприклад, в одному з проєктів для Binance

Напрямки блокчейн-розробки

Часті запитання

Останні роботи

  • image_website-b2b-advance_0.webp
    Розробка сайту компанії B2B ADVANCE
    1452
  • image_web-applications_feedme_466_0.webp
    Розробка веб-додатків для компанії FEEDME
    1310
  • image_websites_belfingroup_462_0.webp
    Розробка веб-сайту для компанії БЕЛФІНГРУП
    1005
  • image_ecommerce_furnoro_435_0.webp
    Розробка інтернет магазину для компанії FURNORO
    1270
  • image_logo-advance_0.webp
    Розробка логотипу компанії B2B Advance
    719
  • image_crm_enviok_479_0.webp
    Розробка веб-додатків для компанії Enviok
    1012

Ми розробляємо пайплайни обробки tick-даних — записи кожної угоди з ціною, об'ємом та стороною. Стандартні OHLCV свічки втрачають мікроструктуру ринку: дисбаланс ліквідності, великі угоди, потік buy/sell. Без якісного пайплайну ML-модель навчається на шумі. Наприклад, в одному з проєктів для Binance (з нашої практики) навантаження сягало 300 000 тиків на секунду — ClickHouse впорався, а PostgreSQL впав на 10 000. Наш 5-річний досвід гарантує надійність. Зв'яжіться з нами — ми готові спроєктувати та реалізувати пайплайн під ваші завдання.

Чому tick-дані важливіші за OHLCV для ML?

При агрегації в 1-хвилинні свічки втрачається до 80% інформації: ви не бачите, як розподілені угоди всередині інтервалу, чи був сплеск об'єму, хто був агресором. Volume bars, dollar bars та imbalance bars зберігають ці сигнали. ML-моделі, навчені на тиках, показують на 15–20% вищу точність у задачах прогнозування напрямку руху ціни.

Проблеми, які вирішуємо

  • Високі навантаження. Біржі генерують до 500 000 тиків на секунду. Стандартні БД не справляються з такою вставкою.
  • Затримки. Для HFT-стратегій latency від отримання тику до сигналу не має перевищувати 10 ms.
  • Зберігання. Tick-дані за рік — це десятки терабайт. Необхідні партиціонування, TTL та ефективне стиснення. Економія на інфраструктурі ClickHouse може сягати 50% порівняно з традиційними реляційними базами.
  • Різноманітність барів. Time bars нерівномірні в періоди низької активності. Volume/dollar/imbalance bars адаптуються до ринкової активності.

Як ми це робимо: стек і кейс

В одному з проєктів для 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 для таких навантажень. Агрегації за місяць виконуються за секунди. Ми гарантуємо, що ваш пайплайн витримає будь-яку ринкову активність.

Як побудувати volume bars з тиків: покроково

  1. Підключіться до WebSocket біржі для отримання агрегованих угод.
  2. Накопичуйте тики в буфері (наприклад, 1000 записів).
  3. При досягненні заданого об'єму закривайте бар і зберігайте його в ClickHouse.
  4. Використовуйте функцію create_volume_bars з прикладу нижче.

Volume bars закриваються при накопиченні заданого об'єму, а не за фіксований часовий інтервал. Це дає рівномірну кількість спостережень незалежно від активності ринку.

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) 

Аналогічно будуються dollar bars (за об'ємом у USD) та imbalance bars (за дисбалансом buy/sell).

Тип бара Критерій закриття Коли використовувати
Time Інтервал часу Висока ліквідність, рівномірна активність
Volume Накопичений об'єм Адаптація до сплесків волатильності
Dollar Накопичений USD-об'єм Інваріантність до ціни активу
Imbalance Дисбаланс buy/sell Пошук розворотних точок

Яку вигоду дає feature engineering з тиків?

З тиків видобуваються ознаки, які покращують якість ML-моделей: потоковий дисбаланс, частота угод, VWAP-відхилення, частка великих угод. У real-time streaming 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-го перцентиля) часто вказують на institutional activity. Аналіз їх спрямованості дає додатковий сигнал.

Як забезпечити latency <10 ms?

Реальна streaming-архітектура:

Binance WebSocket → asyncio consumer → buffer → ClickHouse batch insert → Redis sorted set (last 10k ticks) → Feature calculator (sliding window) → ML inference → Signal output 

Затримка від тику до сигналу — менше 10 ms. Досягається за рахунок асинхронного I/O, буферизації в Redis і попереднього розрахунку фіч на часових вікнах. Економія на кластері ClickHouse порівняно з традиційними БД сягає 50%.

Згідно з документацією ClickHouse, швидкість вставки сягає 500,000 рядків на секунду ClickHouse Documentation.

Процес роботи

Етап Тривалість Результат
Аналітика 1–2 дні Документ з вимогами та схемою даних
Проєктування 2–3 дні Вибір стеку, проєктування схеми БД, визначення типів барів
Реалізація 1–2 тижні Collector, агрегатори, feature engineering, інтеграція з ML-пайплайном
Тестування 3–5 днів Валідація на історичних даних, стрес-тест за швидкістю
Деплой 2–3 дні Розгортання у вашому кластері (Docker/K8s), моніторинг

Строки та що входить

Базова версія (один символ, ClickHouse, Redis) — від 2 тижнів. Повний пайплайн з volume/dollar/imbalance bars, feature engineering та real-time inference — від 4 тижнів. Вартість проєкту варіюється: базова версія — близько 2000 USD, повний пайплайн — від 5000 USD. При цьому економія на зберіганні з ClickHouse може становити до $500 на місяць порівняно з PostgreSQL. Інвестиція в якісний пайплайн окупається за рахунок підвищення точності ML-моделей для торгівлі.

Зазначимо: Що входить в роботу:

  • Архітектурна документація.
  • Вихідний код пайплайну з коментарями.
  • Налаштування ClickHouse, Redis, черг.
  • Інтеграція з вашою ML-інфраструктурою.
  • Навчання команди (2-3 дзвінки).
  • Підтримка 2 місяці після деплою.
Чек-лист для перевірки пайплайну
  • Перевірте швидкість вставки: ClickHouse має вставляти не менше 100K рядків/сек на одному ядрі.
  • Переконайтеся в наявності TTL — без нього диск переповниться за місяць.
  • Налаштуйте моніторинг затримок (latency) кожного етапу.
  • Протестуйте автоматичний реконнект WebSocket при розриві з'єднання.
  • Валідуйте агрегації на історичних даних — порівняйте з еталонними барами.

Типові помилки

  • Занадто дрібне партиціонування (по годинах) — велика кількість партицій деградує продуктивність ClickHouse. Оптимально — по днях.
  • Ігнорування TTL — без автоматичного очищення даних диск переповнюється за місяць.
  • Використання time bars для активів з низькою ліквідністю — більшість свічок буде порожніми.

Розробляємо tick-data пайплайни під ключ понад п'ять років, реалізували 30+ проєктів для криптотрейдингу. Отримайте безкоштовний аналіз ваших даних та рекомендації з оптимізації пайплайну. Зв'яжіться з нами — оцінимо ваш проєкт і запропонуємо оптимальне рішення.