Сбор данных о сделках (trades) с криптобирж
Недавно клиент потерял три дня, пытаясь собрать trades с Binance через REST — пропустил 15% сделок из-за лимитов. При пиковых нагрузках 1200 запросов в минуту не хватало для покрытия 20 торговых пар, и данные приходили с задержкой более секунды. Мы перевели его на WebSocket и получили задержку 2 мс, полнота данных — 99.99%. Такие кейсы — норма: ограничения API, разрывы соединений, разнобой форматов. Сбор trades в реальном времени — инженерная задача, которую мы решаем под ключ. Обрабатываем до 100 000 сделок/сек на одном VPS. Свяжитесь с нами для оценки вашего проекта.
Почему WebSocket — единственный вариант для real-time?
REST polling добавляет 500–2000 мс задержки и пропускает сделки при пиковой нагрузке. WebSocket даёт потоковую передачу с задержкой 1–50 мс. Binance WebSocket documentation рекомендует до 300 стримов на одно соединение. Мы используем несколько соединений для покрытия всех пар.
Управление множеством бирж
CCXT Pro предоставляет единый интерфейс watch_trades для 30+ бирж с автореконнектом. Код ниже подключается к любой CEX за пару строк:
import ccxt.pro as ccxtpro import asyncio async def collect_trades(exchange_id: str, symbols: list[str], queue: asyncio.Queue): exchange = getattr(ccxtpro, exchange_id)({ 'enableRateLimit': True, 'options': {'tradesLimit': 1000}, }) try: while True: try: trades = await exchange.watch_trades_for_symbols(symbols) for trade in trades: await queue.put({ 'exchange': exchange_id, 'symbol': trade['symbol'], 'id': trade['id'], 'price': trade['price'], 'amount': trade['amount'], 'side': trade['side'], 'timestamp': trade['timestamp'], }) except Exception as e: print(f'Error {exchange_id}: {e}, reconnecting...') await asyncio.sleep(1) finally: await exchange.close() CCXT Pro в 10 раз быстрее в разработке, чем собственные WebSocket-коннекторы под каждую биржу. Для масштабирования мы используем кластеризацию: несколько инстансов распределяют нагрузку по разным группам торговых пар. Это позволяет обрабатывать тысячи пар без потери производительности.
Как собирать сделки с DEX?
На DEX сделки — это события смарт-контрактов. Два основных метода: The Graph subgraph — готовые данные через GraphQL. Для Uniswap V3:
{ swaps( first: 100 orderBy: timestamp orderDirection: desc where: { pool: "0x8ad599c3a0ff1de082011efddc58f1908eb6e6d8" } ) { id timestamp amount0 amount1 sqrtPriceX96 tick transaction { id } } } Задержка — 1–5 минут от появления в блоке. Прямой RPC мониторинг — подписка на Swap-события через eth_subscribe. Шаги:
- Подключитесь к RPC узлу через WebSocket.
- Подпишитесь на событие
Swapдля пула. - Декодируйте
sqrtPriceX96в цену. - Обработайте и сохраните данные.
Пример на TypeScript:
import { createPublicClient, webSocket, parseAbiItem } from 'viem'; const SWAP_EVENT = parseAbiItem( 'event Swap(address indexed sender, address indexed recipient, int256 amount0, int256 amount1, uint160 sqrtPriceX96, uint128 liquidity, int24 tick)' ); client.watchContractEvent({ address: UNISWAP_V3_POOL, event: SWAP_EVENT, onLogs: (logs) => { for (const log of logs) { const { amount0, amount1, sqrtPriceX96 } = log.args; const price = sqrtPriceX96ToPrice(sqrtPriceX96, token0Decimals, token1Decimals); processSwap({ price, amount0, amount1, txHash: log.transactionHash }); } } }); Цена в Uniswap V3 хранится как sqrtPriceX96 (Q64.96 fixed point). Декодирование:
function sqrtPriceX96ToPrice(sqrtPriceX96: bigint, d0: number, d1: number): number { const price = Number(sqrtPriceX96 ** 2n * BigInt(10 ** d0)) / Number(BigInt(2 ** 192) * BigInt(10 ** d1)); return price; } Сравнение методов сбора DEX-данных
| Метод | Задержка | Пропускная способность | Сложность реализации |
|---|---|---|---|
| The Graph subgraph | 1–5 мин | Высокая | Низкая |
| Прямой RPC мониторинг | ~500 мс | Средняя | Средняя |
| Парсинг логов через eth_getLogs | ~5 с | Низкая | Высокая |
Сравнение CEX и DEX по сбору trades
| Характеристика | CEX | DEX |
|---|---|---|
| Тип данных | Централизованный API | On-chain события |
| Задержка | 1–50 мс (WebSocket) | 500 мс – 5 мин |
| Надёжность | Высокая (лимиты по IP) | Зависит от ноды |
| Сложность интеграции | Средняя (CCXT) | Высокая (декодирование) |
Как обходить rate limits и блокировки?
Binance: 1200 запросов/мин на IP для REST, WebSocket — до 300 стримов на соединение. Используем несколько соединений по 300 пар.
Bybit и OKX имеют аналогичные ограничения. Bybit отключает WebSocket при отсутствии активности — ping каждые 20 с. Ротация IP работает для REST, но не для WebSocket. Для high-frequency используем несколько VPS в разных датацентрах.
Также настраиваем сжатие и пулы соединений для снижения нагрузки. Опыт нашей команды (50+ интеграций) гарантирует стабильность даже при пиковых объёмах.
Пример конфигурации для работы с rate limits
# Настройка CCXT Pro с контролем лимитов exchange = ccxtpro.binance({ 'enableRateLimit': True, 'rateLimit': 1000, 'options': { 'tradesLimit': 1000, 'watchTrades': {'limit': 100}, }, }) Для кластеризации запускаем несколько таких инстансов, каждый со своим набором символов.
Как нормализовать и хранить trades?
Приводим все сделки к единой схеме с партиционированием по дням:
CREATE TABLE trades ( id BIGSERIAL PRIMARY KEY, exchange VARCHAR(50) NOT NULL, symbol VARCHAR(30) NOT NULL, trade_id VARCHAR(100), price NUMERIC(30, 10) NOT NULL, quantity NUMERIC(30, 10) NOT NULL, side CHAR(4) NOT NULL, ts TIMESTAMPTZ NOT NULL, received_at TIMESTAMPTZ DEFAULT NOW() ) PARTITION BY RANGE (ts); CREATE INDEX ON trades (exchange, symbol, ts DESC); CREATE INDEX ON trades (symbol, ts DESC); TimescaleDB упрощает это через time_bucket для OHLCV-агрегаций. Использование TimescaleDB снижает затраты на хранение по сравнению с PostgreSQL примерно на 30% за счёт сжатия и партиционирования.
Дедупликация
При reconnect WebSocket досылает последние N трейдов. Unique constraint на (exchange, trade_id) предотвращает дубли.
Что входит в работу
Мы поставляем:
- Архитектуру системы сбора (выбор стримов, стыковка CEX + DEX)
- Код на Python/TypeScript с обработкой ошибок и реконнектом
- Нормализацию данных под единую схему
- Деплой на VPS/Kubernetes с мониторингом
- Документацию и обучение команды
- Поддержка 30 дней после релиза
Стоимость разработки системы сбора trades зависит от количества бирж и торговых пар. Для типового набора из 3-5 бирж и 10-20 пар она составляет от 2 000 до 5 000 долларов. Более 10 лет опыта в блокчейн-разработке, 50+ интеграций бирж, сертифицированные инженеры — гарантируем надёжный сбор trades под любую нагрузку. Получите консультацию — мы подготовим архитектуру под ваши объёмы.







