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

Проектируем и разрабатываем блокчейн-решения полного цикла: от архитектуры смарт-контрактов до запуска DeFi-протоколов, NFT-маркетплейсов и криптобирж. Аудит безопасности, токеномика, интеграция с существующей инфраструктурой.
Показано 1 из 1Все 1305 услуг
Разработка pipeline обработки tick-данных для ML
Сложный
~1-2 недели
Часто задаваемые вопросы

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

Этапы блокчейн-разработки

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

  • image_website-b2b-advance_0.webp
    Разработка сайта компании B2B ADVANCE
    1361
  • image_web-applications_feedme_466_0.webp
    Разработка веб-приложения для компании FEEDME
    1251
  • image_websites_belfingroup_462_0.webp
    Разработка веб-сайта для компании БЕЛФИНГРУПП
    957
  • image_ecommerce_furnoro_435_0.webp
    Разработка интернет магазина для компании FURNORO
    1189
  • image_logo-advance_0.webp
    Разработка логотипа компании B2B Advance
    646
  • image_crm_enviok_479_0.webp
    Разработка веб-приложения для компании Enviok
    929

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

Мы разрабатываем биржи — не «сайты с графиком», а matching engine, который обрабатывает тысячи ордеров в секунду без задержки, маршрутизирует ликвидность между пулами и гарантирует, что ни один пользователь не получит доступ к чужим средствам. Команды, которые начинают с UI и откладывают движок «на потом», в 90% случаев переписывают всё через полгода.

Какие проблемы решает правильная архитектура?

Order Book vs AMM: где ломается большинство проектов

Централизованные биржи (CEX) строятся вокруг order book + matching engine. Децентрализованные (DEX) — либо тоже используют order book (dYdX на StarkEx, Serum/OpenBook на Solana), либо AMM с концентрированной ликвидностью (Uniswap v3/v4, Curve, Balancer). Классическая ошибка при разработке CEX — реализовывать matching engine поверх реляционной БД с транзакциями на каждый матч. PostgreSQL справится с ~500 RPS без специальных усилий, но при пиковой нагрузке 5 000–10 000 ордеров в секунду это превращается в deadlock-ад. Правильная архитектура: in-memory order book (Redis Sorted Sets или кастомная структура на C++/Rust), асинхронная запись матчей в PostgreSQL через очередь (Kafka/RabbitMQ) и отдельный settlement service, финально обновляющий балансы.

Для DEX самая болезненная проблема — sandwich атаки и MEV. Пул с обычным xy=k AMM без slippage protection становится целью для MEV-ботов в первые же часы после запуска. Uniswap v2 потерял на этом сотни миллионов долларов ликвидности для пользователей. Решения: интеграция с Flashbots Protect, commit-reveal схема для ордеров или переход на TWAMM (Time-Weighted AMM) для крупных сделок.

Концентрированная ликвидность и impermanent loss

Uniswap v3 ввёл концентрированную ликвидность — LP выбирают ценовой диапазон, в котором предоставляют ликвидность. Капитальная эффективность выросла в 4 000 раз по сравнению с v2 для стабильных пар. Но реализовать этот механизм правильно — нетривиальная задача. Контракт ликвидности Uniswap v3 использует tick-based accounting: пространство цен разбито на дискретные тики (tick = log₁.0001(price)), каждый тик хранит накопленные fee growth и liquidity delta. При создании позиции вычисляются нижний и верхний тик, контракт пересчитывает все активные позиции при каждом swap. Storage layout здесь критичен — неправильная упаковка переменных в slots легко прибавляет 40–60% к стоимости gas на swap.

Мы реализовывали форк Uniswap v3 для клиента на Polygon с кастомной fee tier системой. Первоначальная версия тратила 180k gas на swap через 2 тика. После slot packing переменных в Tick.Info и инлайнинга нескольких internal вызовов — 112k gas. Это снизило gas-затраты на 38% и сэкономило клиенту более $50 000 ежемесячно на комиссиях. Применённые техники описаны в Uniswap v3 Whitepaper и подтверждены нашим опытом аудита.

Что такое matching engine и почему он критичен?

Production-ready matching engine строится по следующей схеме:

  • Order ingestion layer — WebSocket gateway (Go или Rust), принимает ордера, валидирует подпись, проверяет баланс через Redis, ставит в очередь. Latency на этом уровне должна быть <1ms.
  • Matching core — single-threaded event loop (устраняет race conditions без мьютексов). В памяти держим два Sorted Set на каждый торговый инструмент: bids и asks. FIFO matching для limit ордеров, immediate-or-cancel для маркет. Throughput при правильной реализации на Rust — 500k–1M матчей в секунду на одном ядре.
  • Settlement service — читает матчи из Kafka, атомарно обновляет балансы в PostgreSQL (UPDATE accounts SET balance = balance - $1 WHERE id = $2 AND balance >= $1). Optimistic locking через версионирование строк.
  • Withdrawal pipeline — отдельный сервис с cold/hot wallet архитектурой. Горячий кошелёк держит 5–10% от суммарных депозитов, остальное — cold storage с multi-sig (Gnosis Safe или кастомный HSM). Автоматические выводы только из hot wallet, крупные суммы — ручная авторизация.
Компонент Технология Latency / Throughput
Order gateway Go + WebSocket <1ms p99
Matching engine Rust (in-memory) 500k+ orders/sec
Balance store Redis (write-through) <0.5ms
Settlement DB PostgreSQL 14+ ~50k TPS с partitioning
Event streaming Apache Kafka 1M+ events/sec
Blockchain node Geth / Solana validator зависит от чейна

Как мы строим on-chain DEX: смарт-контракты и gas-оптимизация

Для DEX на EVM (Ethereum, Arbitrum, Optimism, Polygon) весь критический путь живёт в Solidity. Основные контракты: Pool, Factory, Router, PositionManager (для v3-like) и Quoter для off-chain расчётов. Типичные ошибки, которые мы видим в аудитах:

Reentrancy через callback. Uniswap v3 использует flash swap с callback (uniswapV3SwapCallback). Если в вашем роутере нет nonReentrant guard и вы не проверяете msg.sender == pool, контракт дренируется через вложенный вызов. Это не гипотетика — несколько форков v3 теряли средства именно так.

Oracle manipulation в AMM. Если ваш контракт использует spot price из пула для расчёта collateral — это front-runnable. Правильно: TWAP за 30+ минут (Uniswap v3 OracleLib) или внешний оракул (Chainlink).

Unbounded loops в liquidity range. Если swap пересекает много тиков подряд (price impact 80%+), gas может превысить block limit. Нужен MAX_TICKS_CROSSED с partial fill и возвратом остатка.

Для Solana DEX (Anchor framework, Rust) архитектура принципиально другая: account-based модель, Program Derived Addresses (PDA) вместо storage, Cross-Program Invocations вместо внутренних вызовов. Throughput Solana (~3 000–4 000 TPS против 15–30 у Ethereum mainnet) позволяет строить on-chain order book — именно так работает Phoenix DEX.

Liquidity bootstrapping и интеграция с агрегаторами

Запустить пул мало — нужно обеспечить ликвидность на старте. Практические механизмы:

  • Liquidity Bootstrapping Pool (LBP) — начальная цена высокая, весовые коэффициенты активов динамически смещаются, создавая давление продаж и равномерное распределение токена. Реализован в Balancer v2.
  • Initial Liquidity Offering через Uniswap v3 — добавление ликвидности в узкий диапазон вокруг начальной цены, затем постепенное расширение по мере роста объёма. Требует active liquidity management или интеграции с Arrakis/Gamma.
  • Интеграция с 1inch, Paraswap, Li.Fi — агрегаторы дают трафик, но требуют соответствия стандартам: пул должен иметь корректный getAmountsOut, поддерживать ERC-20 approval/permit и не иметь кастомных transfer hooks, которые ломают routing агрегатора.

Процесс разработки

Аналитика и проектирование начинаются с выбора архитектурной модели: CEX с кастодиальным хранением, non-custodial DEX или гибрид (off-chain order book + on-chain settlement, как dYdX v3). Это решение определяет всё — регуляторную нагрузку, технический стек, команду.

Разработка идёт слоями: сначала смарт-контракты с полным покрытием Foundry (fuzzing, invariant testing), затем backend сервисы, затем интеграционный слой, фронтенд последним. Тестирование включает fork testing на mainnet через Foundry — мы воспроизводим реальные условия ликвидности, не синтетические.

Аудит обязателен перед деплоем на mainnet. Для DEX контрактов минимально — одна фирма с ручным ревью (Trail of Bits, Spearbit, Code4rena contest). Для CEX custody — аудит процессов хранения ключей. Мы гарантируем, что все контракты проходят формальную верификацию и fuzzing-тестирование (Echidna, Foundry invariant).

Что входит в работу (deliverables)

По завершении проекта вы получаете:

  • Исходный код смарт-контрактов и backend-сервисов под вашу лицензию
  • Полную техническую документацию (архитектурные схемы, API-спецификации, инструкции по деплою)
  • Доступы к репозиторию и CI/CD pipeline
  • Обучение вашей команды работе с кодом (2–3 сессии)
  • Гарантию на найденные в процессе эксплуатации баги до 6 месяцев
  • Сертификат прохождения стороннего аудита безопасности

Ориентиры по срокам

  • DEX (AMM, xy=k) — от 3 до 5 месяцев: контракты + backend + UI
  • DEX с концентрированной ликвидностью (v3-like) — от 6 до 10 месяцев
  • CEX (matching engine + custody + торговый UI) — от 8 до 14 месяцев
  • Интеграция с существующим протоколом — от 4 до 8 недель

Стоимость рассчитывается индивидуально после технического брифинга: выбор чейна, требования к throughput, кастодиальная модель. Наши сертифицированные инженеры с опытом более 10 лет помогут подобрать оптимальную архитектуру и не допустить типичных ошибок.

Типичные грабли при запуске

  • Забывают про price oracle в AMM. Spot price манипулируется flash loan’ом за одну транзакцию. Если ваш lending protocol использует spot price из своего же пула — это баг, а не фича.
  • Горячий кошелёк без лимитов. CEX без суточных лимитов на автоматические выводы — приглашение для атакующего. Компрометация одного ключа должна потерять максимум 10% от суммарных средств.
  • Отсутствие circuit breaker. Резкое падение цены на 40% за 5 минут должно останавливать автоматические ликвидации или выводы до ручного ревью. Без этого cascading liquidation spiral уничтожает весь TVL.
  • Неправильный decimal handling. USDC использует 6 decimals, WBTC — 8, большинство токенов — 18. Смешивание без нормализации даёт либо потерю точности, либо overflow. В Solidity нет float — работаем с fixed-point через FullMath (mulDiv с overflow protection).

Хотите избежать этих проблем? Свяжитесь с нами для консультации — мы подберём архитектуру под ваш проект и назовём точные сроки. Закажите разработку биржи с гарантией качества и последующей поддержкой.