Парсинг крипто-даних — лише перший крок. Коли дані надходять із п'яти бірж, трьох блокчейн-мереж та двох соціальних платформ — кожне джерело присилає їх у своєму форматі. 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). - Числові формати:
nullvs"0"vs0vs відсутність поля — для нульових обсягів;-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 до нормалізації. Якщо виявлена помилка в логіці нормалізації — можна перепрогнати за вихідними даними без повторного збору.
Як впровадити нормалізаційний шар: покроковий процес
- Аналіз джерел: визначаємо всі джерела даних (біржі, блокчейни, API), збираємо зразки форматів.
- Проектування схем: створюємо Pydantic/Zod схеми для кожного джерела з версіонуванням.
- Розробка трансформацій: пишемо функції нормалізації для кожного поля (timestamp, суми, символи).
- Тестування та моніторинг: проганяємо на історичних даних, налаштовуємо алерти.
Що входить у роботу
При замовленні під ключ ви отримуєте:
- Готовий нормалізаційний шар для ваших джерел (до 7 у базовому варіанті)
- Документацію схем та API
- Доступ до репозиторію з кодом та тестами
- Навчання команди роботі з системою
- Підтримку протягом 1 місяця після запуску
Замовте розробку нормалізаційного шару. Зв'яжіться з нами, щоб обговорити ваш проект. Ми гарантуємо прозорий процес та індивідуальний підхід.







