Розробка pipeline обробки order book даних для ML

Проєктуємо та розробляємо блокчейн-рішення повного циклу: від архітектури смарт-контрактів до запуску DeFi-протоколів, NFT-маркетплейсів та криптобірж. Аудит безпеки, токеноміка, інтеграція з наявною інфраструктурою.
Показано 1 з 1Усі 1305 послуг
Розробка pipeline обробки order book даних для 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

Розробка pipeline обробки order book даних для ML

Повний стакан ордерів містить всю інформацію про ліквідність на біржі. Але зібрати його, нормалізувати і перетворити на ознаки для машинного навчання — нетривіальне інженерне завдання. Ми побудували production-grade pipeline для Binance, Bybit та OKX, який обробляє до 10 000 оновлень за секунду. Наш досвід включає інтеграцію з 15+ криптобіржами та зберігання близько 5 ТБ даних на місяць. Ми маємо 5+ років досвіду в криптоінфраструктурі та 50+ реалізованих проектів під ключ. Повний L2 стакан описує кожен рівень ціни з обсягом — це основа для побудови короткострокових прогнозів. Гарантуємо стабільний збір при пікових навантаженнях та консистентність знімків.

Замовники часто приходять з сирими WebSocket-стрімами, не знаючи, як синхронізувати diff stream зі REST-знімком. Помилка в one-off призводить до роз'їзду стакана і невірних сигналів. Ми вирішуємо цю проблему на рівні архітектури колектора.

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

  • Обсяг даних. Повний L2 стакан на Binance містить 5000 рівнів з кожного боку. При оновленнях кожні 100 мс це генерує десятки гігабайт на день. Наївне зберігання в PostgreSQL вб'є продуктивність.
  • Гонка станів. WebSocket diff stream приходить асинхронно. Без синхронізації зі REST-знімком стакан роз'їжджається — ціна йде в неіснуючі рівні.
  • Формат даних. Кожна біржа віддає стакан по-своєму: Binance — вкладені масиви, Coinbase — JSON з різними ключами. Потрібен єдиний інтерфейс.

Як синхронізувати WebSocket diff stream зі REST-знімком?

Алгоритм простий: відкриваємо WebSocket, отримуємо перший diff stream, одразу запитуємо REST-знімок з повним станом. Далі кожне оновлення накладаємо на локальний стакан. Для контролю використовуємо lastUpdateId: застосовуємо лише повідомлення з u > lastUpdateId. Якщо послідовність порушена — перезапитуємо знімок. Цей підхід виключає роз'їзд стакана навіть при високій волатильності.

Як зібрати order book через WebSocket: покроковий алгоритм

  1. Встановлення з'єднання: через wss://stream.binance.com:9443/ws/btcusdt@depth@100ms (аналог для інших бірж).
  2. Первісний REST-знімок: синхронізація через updateId для забезпечення консистентності.
  3. Інкрементальні оновлення: кожне повідомлення diff stream накладається на поточний стан стакана.
  4. Збереження знімків: із заданою періодичністю (кожне N-те оновлення) фіксується повний стан для подальшого feature engineering.
Приклад коду колектора
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 — оптимальне сховище для order book?

Повний L2 стакан — величезний обсяг. ClickHouse в 10 разів швидше PostgreSQL на колонкових агрегаціях. Згідно документації ClickHouse, колонкова СУБД забезпечує стиснення до 10 разів та швидкість запису більш ніж 1 млн рядків за секунду. Порівняйте:

СУБД Швидкість запису (рядків/с) Стиснення Агрегації за часом
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 ТБ даних. Для великих проектів економія може досягати $30,000 на рік. Вартість розробки базового pipeline під ключ — від $5,000 до $15,000.

Feature engineering з order book

На основі зібраних знімків будуємо ознаки. Базові: OBI (order book imbalance), спред, глибина. Додаткові: ковзні середні 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

Як оцінити якість прогнозу mid-price?

Для короткострокового прогнозу mid-price (через N оновлень стакана) використовуємо метрики accuracy, precision та 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]

Типові помилки при розробці order book pipeline

Навіть досвідчені команди допускають помилки: ігнорування перекосів стакана в моменти високої волатильності, неправильна обробка подій lastUpdateId, відсутність перевірки консистентності після реконекту. Ми стикалися з проектом, де через пропущені диффи стакан розійшовся на 20% — модель показувала хибні сигнали. Рішення — вбудовування перевірок контрольних сум та автоматичне відновлення повного знімку при виявленні невідповідності.

Що входить у розробку pipeline під ключ

  • Вихідний код колектора та пайплайна (асинхронний Python).
  • Дампи тестових даних для offline-тестування.
  • README з докладними прикладами використання.
  • Міграції схеми ClickHouse з TTL.
  • Навчання вашої команди роботі з pipeline.

Етапи роботи та строки

Етап Тривалість Результат
Аналітика 2-3 дні Специфікація API, обсягів
Проектування 2-3 дні Схема зберігання, вибір ознак
Реалізація 5-10 днів Колектор, пайплайн, код
Тестування 3-5 днів Симуляція 24h, звіти
Деплой 2-3 дні Docker, моніторинг

Базовий pipeline для однієї біржі з моделлю LightGBM — від 14 до 30 робочих днів. Точну оцінку даємо після безкоштовного аудиту ваших даних. Замовте аналіз — і ми підберемо оптимальну архітектуру під ваш обсяг стакана. Пишіть нам для консультації — оцінимо проєкт безкоштовно.

Розробка бірж: matching engine визначає успіх

Ми розробляємо біржі, де matching engine обробляє тисячі ордерів на секунду без затримки, маршрутизує ліквідність між пулами та гарантує, що жоден користувач не отримає доступ до чужих коштів. Команди, які починають з UI і відкладають движок «на потім», у 90% випадків переписують все через півроку. Наш досвід — 15+ запущених біржових проєктів. Оцініть ваш проєкт — отримайте консультацію.

Типові проблеми архітектури бірж

Order Book vs AMM

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

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

Як захистити DEX від MEV-атак?

Flashbots Protect дозволяє відправляти транзакції напряму в блок без публічного mempool. Commit-reveal схема робить неможливим front-running, приховуючи параметри ордера до моменту виконання. Для децентралізованих order book-бірж (на кшталт dYdX) це критично — без захисту MEV-боти викачують прибуток маркет-мейкерів. Ми реалізовували таку інтеграцію для клієнта на Arbitrum: після підключення Flashbots частка sandwich-атак знизилась з 12% до 0.2% від усіх угод.

Концентрована ліквідність та 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% і зекономило клієнту понад $5,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: смарт-контракти та газ-оптимізація

Для 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 і поверненням залишку.

Тип помилки Наслідок Рішення
Reentrancy Втрата коштів через вкладений виклик nonReentrant guard + перевірка caller
Oracle manipulation Маніпуляція ціною через flash loan TWAP або зовнішній оракул
Unbounded loops Транзакція не влазить у блок Partial fill + ліміт тиків

Як оптимізувати газ для смарт-контрактів DEX?

Оптимізація gas включає packing змінних у storage slots, використання inline assembly для критичних операцій та мінімізацію зовнішніх викликів. Правильне розміщення полів у структурі Tick.Info дозволяє зменшити gas на 20–30% порівняно з базовою реалізацією. Для 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 агрегатора.

Використовуйте LBP для створення початкового цінового діапазону, а потім підключайте агрегатори для забезпечення постійного потоку замовлень. Активне управління ліквідністю через професійні протоколи допомагає уникнути втрат від impermanent loss. Наш досвід — 15+ запущених біржових проєктів, які пройшли незалежний аудит. Середня економія клієнтів на gas-комісіях після оптимізації — $5,000 щомісяця.

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

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

Як проходить тестування смарт-контрактів?

Ми використовуємо Foundry для unit-тестів, fuzzing та invariant testing. Fork testing на mainnet дозволяє відтворити реальні умови ліквідності, що критично для верифікації поведінки контрактів.

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

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

Що входить в роботу (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).

Зв'яжіться з нами для консультації — ми підберемо архітектуру під ваш проєкт і назвемо точні терміни. Замовте розробку біржі з гарантією якості та подальшою підтримкою.