Разработка системы нормализации данных из крипто-источников

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

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

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

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

  • image_website-b2b-advance_0.webp
    Разработка сайта компании B2B ADVANCE
    1452
  • 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 месяца после запуска

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