Розробка системи нормалізації даних із крипто-джерел

Парсинг крипто-даних — лише перший крок. Коли дані надходять із п'яти бірж, трьох блокчейн-мереж та двох соціальних платформ — кожне джерело присилає їх у своєму форматі. Binance повертає timestamps у мілісекундах, OKX — у секундах, Telegram — у UTC datetime, on-chain дані — у Unix секундах з блоку.

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

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

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

  • image_website-b2b-advance_0.webp
    Розробка сайту компанії B2B ADVANCE
    1451
  • image_web-applications_feedme_466_0.webp
    Розробка веб-додатків для компанії FEEDME
    1309
  • 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
    1011

Парсинг крипто-даних — лише перший крок. Коли дані надходять із п'яти бірж, трьох блокчейн-мереж та двох соціальних платформ — кожне джерело присилає їх у своєму форматі. Binance повертає timestamps у мілісекундах, OKX — у секундах, Telegram — у UTC datetime, on-chain дані — у Unix секундах з блоку. Суми всюди різні: десь wei, десь Gwei, десь string з плаваючою точкою. Ми будуємо нормалізаційний шар, який перетворює цей хаос на єдиний, передбачуваний формат. Оцінимо ваш проект за 1 день — просто зв'яжіться з нами.

Як нормалізація даних впливає на надійність DeFi-систем

Помилка в одному тікері або втрата точності на шостому знаку може призвести до втрати коштів або невірних метрик. Наш досвід — 10+ років у блокчейн-розробці — показує, що 80% інцидентів з даними пов'язані саме з неправильною нормалізацією. Без неї жоден ковзний хедж або арбітраж не працює. Порівняємо: нормалізований pipeline обробляє дані в 3 рази швидше за ad-hoc скрипти, а ймовірність помилки знижується на порядок.

Проблеми гетерогенних даних

Перелічимо конкретні розбіжності, які зустрічаються в реальних проектах:

  • Часові мітки: Unix milliseconds (Binance, most CEX), Unix seconds (Ethereum blocks, Chainlink), ISO 8601 strings (деякі REST API), Relative ("2 hours ago") — у social data scraping, Timezone-aware vs naive datetimes.
  • Суми та ціни: Wei (10^-18 ETH) — on-chain Ethereum, Lamports (10^-9 SOL) — on-chain Solana, String з decimals ("1234.567890") — Binance REST, Integer з fixed decimals (100000000 = 1 BTC у деяких бірж), Float64 — втрата точності на великих числах.
  • Ідентифікатори активів: BTCUSDT (Binance), BTC-USDT (OKX), BTC/USDT (ccxt standard), tBTCUST (Bitfinex), ERC-20 address (0x2260fac...) vs ticker (WBTC), CoinGecko ID ("bitcoin") vs CMC ID (1).
  • Числові формати: null vs "0" vs 0 vs відсутність поля — для нульових обсягів; -0.0 — валідне значення в Python/JS float, неочевидна поведінка при порівнянні; NaN — іноді зустрічається в JSON від сторонніх API.

Як побудувати нормалізаційний шар?

Система складається з трьох шарів:

Raw Data (from scrapers) ↓ [Validation Layer] — відкидаємо невалідні записи, логуємо помилки ↓ [Transformation Layer] — приводимо до єдиного формату ↓ [Enrichment Layer] — додаємо derived поля (USD-вартість, нормалізований тікер) ↓ Normalized Storage 

Validation Layer

Перед трансформацією — явна валідація вхідних даних. Використовуємо Pydantic v2 для Python:

from pydantic import BaseModel, field_validator, model_validator from decimal import Decimal from datetime import datetime from typing import Optional class RawTradeEvent(BaseModel): """Схема для сирих trade подій від будь-якої біржі""" exchange: str raw_symbol: str raw_price: str | float | int raw_quantity: str | float | int raw_timestamp: int | str | float side: str # 'buy'/'sell' або 'BUY'/'SELL' або 1/2 raw_trade_id: str | int @field_validator('raw_price', 'raw_quantity', mode='before') @classmethod def coerce_to_string(cls, v): if isinstance(v, float): return f"{v:.10f}" return str(v) @field_validator('side', mode='before') @classmethod def normalize_side(cls, v): s = str(v).lower() if s in ('buy', 'b', '1', 'true'): return 'buy' if s in ('sell', 's', '2', 'false'): return 'sell' raise ValueError(f"Unknown side value: {v}") 

Невалідні записи не обрушують весь pipeline — вони логуються в окрему таблицю validation_errors з raw-контекстом і причиною помилки.

Transformation Layer

Приведення до канонічного формату:

from dataclasses import dataclass from decimal import Decimal, ROUND_DOWN from datetime import datetime, timezone @dataclass class NormalizedTrade: exchange: str symbol: str # canonical: "BTC/USDT" price: Decimal # завжди Decimal, жодних float quantity: Decimal quote_quantity: Decimal # price * quantity side: str # 'buy' або 'sell' timestamp: datetime # UTC timezone-aware trade_id: str # рядок, унікальний в межах біржі def normalize_trade(raw: RawTradeEvent) -> NormalizedTrade: return NormalizedTrade( exchange=raw.exchange, symbol=normalize_symbol(raw.raw_symbol, raw.exchange), price=parse_decimal(raw.raw_price), quantity=parse_decimal(raw.raw_quantity), quote_quantity=parse_decimal(raw.raw_price) * parse_decimal(raw.raw_quantity), side=raw.side, timestamp=normalize_timestamp(raw.raw_timestamp), trade_id=str(raw.raw_trade_id), ) def normalize_timestamp(raw: int | str | float) -> datetime: """Приводить будь-який timestamp до UTC datetime""" if isinstance(raw, str): dt = datetime.fromisoformat(raw.replace('Z', '+00:00')) return dt.astimezone(timezone.utc) ts = float(raw) if ts > 1e12: ts = ts / 1000 return datetime.fromtimestamp(ts, tz=timezone.utc) def parse_decimal(value: str) -> Decimal: """Безпечна конвертація в Decimal""" try: d = Decimal(str(value)) if d.is_nan() or d.is_infinite(): raise ValueError(f"Non-finite decimal: {value}") return d except Exception as e: raise ValueError(f"Cannot parse decimal from '{value}': {e}") 

У Python Decimal забезпечує точне зберігання чисел з плаваючою точкою.

Symbol normalization

Мапінг тікерів між біржами — окреме завдання. Використовуємо ccxt-сумісний формат BASE/QUOTE:

SYMBOL_MAPPINGS = { "binance": { "BTCUSDT": "BTC/USDT", "ETHUSDT": "ETH/USDT", }, "okx": { "BTC-USDT": "BTC/USDT", "BTC-USDT-SWAP": "BTC/USDT:USDT", # perpetual }, "bybit": { "BTCUSDT": "BTC/USDT", "BTCPERP": "BTC/USDT:USDT", }, } def normalize_symbol(raw_symbol: str, exchange: str) -> str: exchange_map = SYMBOL_MAPPINGS.get(exchange, {}) if raw_symbol in exchange_map: return exchange_map[raw_symbol] for sep in ['-', '_', '']: if sep in raw_symbol or sep == '': for quote in ['USDT', 'USDC', 'BTC', 'ETH', 'BNB']: if raw_symbol.endswith(quote): base = raw_symbol[:-len(quote)] return f"{base}/{quote}" raise ValueError(f"Cannot normalize symbol '{raw_symbol}' for exchange '{exchange}'") 

Чому важливий schema registry?

Джерела даних змінюються. Binance оновив API — додалося поле, змінився формат timestamp. Без версіонування схем зламається вся нормалізація. Schema registry (аналог Confluent Schema Registry для Kafka) вирішує це: кожен запис містить версію схеми джерела, старі дані не ламаються, а нормалізацію можна перепрогнати при виправленні логіки без повторного збору.

SCHEMA_VERSIONS = { "binance_trade": { "v1": BinanceTradeV1Schema, # попередня версія API "v2": BinanceTradeV2Schema, # після оновлення: додано quoteQty } } def get_schema(source: str, version: str): return SCHEMA_VERSIONS[source][version] 

Моніторинг якості даних

Нормалізація без моніторингу — це ілюзія якості. Ключові метрики:

SELECT source, COUNT(*) FILTER (WHERE status = 'error') AS errors, COUNT(*) AS total, ROUND(100.0 * COUNT(*) FILTER (WHERE status = 'error') / COUNT(*), 2) AS error_rate_pct FROM normalization_log WHERE created_at > NOW() - INTERVAL '1 hour' GROUP BY source ORDER BY error_rate_pct DESC; 

Алерт при error_rate > 5% для будь-якого джерела — значить змінився формат даних і потрібно оновити схему. Cross-source consistency check: одна і та ж ціна BTC в один час не повинна розходитися між біржами більш ніж на 0.5%.

Метрики якості нормалізації:

Метрика Опис Поріг алерту
Error rate Частка невалідних записів >5%
Cross-source diff Розбіжність ціни BTC між біржами >0.5%
Latency Затримка від scrap до нормалізації >10 сек

Технологічний стек

Компонент Вибір
Валідація схем Pydantic v2 (Python) або Zod (TypeScript)
Обробка числових значень Python decimal.Decimal, PostgreSQL numeric
Черга Redis Streams або Kafka
Зберігання PostgreSQL (normalized) + raw backup у S3
Schema registry Custom або Confluent Schema Registry
Моніторинг якості dbt tests + Prometheus метрики

Сирі дані завжди зберігаємо в S3 до нормалізації. Якщо виявлена помилка в логіці нормалізації — можна перепрогнати за вихідними даними без повторного збору.

Як впровадити нормалізаційний шар: покроковий процес

  1. Аналіз джерел: визначаємо всі джерела даних (біржі, блокчейни, API), збираємо зразки форматів.
  2. Проектування схем: створюємо Pydantic/Zod схеми для кожного джерела з версіонуванням.
  3. Розробка трансформацій: пишемо функції нормалізації для кожного поля (timestamp, суми, символи).
  4. Тестування та моніторинг: проганяємо на історичних даних, налаштовуємо алерти.

Що входить у роботу

При замовленні під ключ ви отримуєте:

  • Готовий нормалізаційний шар для ваших джерел (до 7 у базовому варіанті)
  • Документацію схем та API
  • Доступ до репозиторію з кодом та тестами
  • Навчання команди роботі з системою
  • Підтримку протягом 1 місяця після запуску

Замовте розробку нормалізаційного шару. Зв'яжіться з нами, щоб обговорити ваш проект. Ми гарантуємо прозорий процес та індивідуальний підхід.