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

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

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

Часто задаваемые вопросы

Последние работы

  • 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+ проектов для криптотрейдинга. Получите бесплатный анализ ваших данных и рекомендации по оптимизации пайплайна. Свяжитесь с нами — оценим ваш проект и предложим оптимальное решение.