asyncpg.Исходный код
Для обеспечения максимальной пропускной способности и исключения блокировок основного потока мы используем асинхронный стек: asyncio для управления событиями, websockets для получения данных от Binance и asyncpg для эффективной работы с PostgreSQL/TimescaleDB. Запись тиков производится пачками (батчами) по достижении лимита размера или таймаута, что минимизирует накладные расходы на транзакции.
import asyncio
import json
import logging
from datetime import datetime, timezone
import asyncpg
import websockets
# Настройка логирования для мониторинга производительности
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(levelname)s - %(message)s'
)
# Конфигурация подключения
DB_DSN = 'postgres://postgres:password@localhost:5432/ticks_db'
BINANCE_WS_URL = 'wss://stream.binance.com:9443/ws/btcusdt@trade'
BATCH_SIZE = 1000 # Оптимальный размер пачки для HFT
BATCH_TIMEOUT = 1.0 # Максимальное время ожидания записи пачки (сек)
async def init_db(pool):
'''Инициализация таблицы и создание гипертаблицы в TimescaleDB'''
async with pool.acquire() as conn:
# Создаем стандартную таблицу
await conn.execute('''
CREATE TABLE IF NOT EXISTS binance_ticks (
timestamp TIMESTAMPTZ NOT NULL,
symbol VARCHAR(10) NOT NULL,
price NUMERIC NOT NULL,
quantity NUMERIC NOT NULL,
is_buyer_maker BOOLEAN NOT NULL
);
''')
# Преобразуем таблицу в гипертаблицу TimescaleDB
try:
await conn.execute('''
SELECT create_hypertable(
'binance_ticks',
'timestamp',
if_not_exists => TRUE,
migrate_data => TRUE
);
''')
logging.info('TimescaleDB hypertable успешно инициализирована.')
except Exception as e:
logging.warning(f'Замечание при создании гипертаблицы: {e}')
async def save_batch(pool, batch):
'''Асинхронная пакетная вставка данных в БД'''
if not batch:
return
async with pool.acquire() as conn:
try:
# Использование executemany для высокопроизводительной вставки
await conn.executemany('''
INSERT INTO binance_ticks (timestamp, symbol, price, quantity, is_buyer_maker)
VALUES ($1, $2, $3, $4, $5)
''', batch)
logging.info(f'Успешно записан пакет из {len(batch)} тиков.')
except Exception as e:
logging.error(f'Критическая ошибка при записи пакета в БД: {e}')
async def stream_ticks():
'''Основной цикл стриминга и накопления батча'''
pool = await asyncpg.create_pool(DB_DSN, min_size=5, max_size=20)
await init_db(pool)
batch = []
last_flush = asyncio.get_event_loop().time()
async def flush_periodically():
'''Фоновая задача для сброса буфера по таймауту'''
nonlocal last_flush
while True:
await asyncio.sleep(0.1)
now = asyncio.get_event_loop().time()
if batch and (now - last_flush >= BATCH_TIMEOUT):
to_write = list(batch)
batch.clear()
last_flush = now
asyncio.create_task(save_batch(pool, to_write))
# Запуск фонового таймера сброса данных
flush_task = asyncio.create_task(flush_periodically())
try:
async with websockets.connect(BINANCE_WS_URL) as ws:
logging.info('Успешное подключение к Binance WebSocket.')
while True:
try:
message = await ws.recv()
data = json.loads(message)
# Парсинг полей тика
ts = datetime.fromtimestamp(data['E'] / 1000.0, tz=timezone.utc)
symbol = data['s']
price = float(data['p'])
qty = float(data['q'])
is_buyer_maker = data['m']
batch.append((ts, symbol, price, qty, is_buyer_maker))
# Запись при достижении лимита размера батча
if len(batch) >= BATCH_SIZE:
to_write = list(batch)
batch.clear()
last_flush = asyncio.get_event_loop().time()
asyncio.create_task(save_batch(pool, to_write))
except websockets.ConnectionClosed:
logging.warning('WebSocket соединение закрыто. Попытка переподключения...')
break
except Exception as e:
logging.error(f'Ошибка обработки тика: {e}')
finally:
flush_task.cancel()
if batch:
await save_batch(pool, batch)
await pool.close()
if __name__ == '__main__':
try:
asyncio.run(stream_ticks())
except KeyboardInterrupt:
logging.info('Стриминг остановлен пользователем.')
Разбор параметров
DB_DSN: Строка подключения к базе данных PostgreSQL/TimescaleDB. Содержит имя пользователя, пароль, хост, порт и целевую базу данных.BINANCE_WS_URL: Адрес WebSocket-эндпоинта Binance для получения сделок в реальном времени (@tradestream).BATCH_SIZE: Размер буфера накопления тиков. Значение1000является оптимальным компромиссом между потреблением оперативной памяти и снижением нагрузки на дисковую подсистему СУБД.BATCH_TIMEOUT: Временной интервал (в секундах), по истечении которого накопленные в буфере тики принудительно записываются в базу данных, даже если лимитBATCH_SIZEне был достигнут. Это предотвращает задержку актуализации данных при низкой активности рынка.create_hypertable: Встроенная функция TimescaleDB, которая преобразует стандартную таблицу в гипертаблицу. Она автоматически разбивает данные на партиции (чанки) по временному интервалу (по умолчанию — 7 дней), обеспечивая стабильную скорость записи и чтения при терабайтных объемах.
Как запустить
Перед запуском скрипта убедитесь, что у вас развернут экземпляр TimescaleDB. Проще всего запустить его через Docker:
docker run -d --name timescaledb -p 5432:5432 -e POSTGRES_PASSWORD=password timescale/timescaledb:latest-pg15 После этого установите необходимые Python-зависимости:
pip install asyncpg websockets Создайте базу данных ticks_db через любой SQL-клиент или консоль psql:
CREATE DATABASE ticks_db; Запустите скрипт. Он автоматически создаст таблицу binance_ticks, преобразует её в гипертаблицу и начнет потоковое сохранение маркет-даты.
Накопленные исторические данные можно использовать для тестирования торговых стратегий или для динамического управления активами. Например, вы можете интегрировать полученные тики в Скрипт автоматического ребалансирования криптовалютного портфеля на Python через CCXT для точного расчета весов портфеля в реальном времени.
Если вы решите расширить архитектуру для сбора приватных данных (балансы, ордера) или работы с другими биржами, вам потребуется реализовать безопасную аутентификацию. Для этого изучите руководство: Генерация подписи HMAC SHA256 для Bybit API V5 на Python.
При масштабировании системы сбора данных на десятки торговых пар вы неизбежно столкнетесь с жесткими лимитами бирж по IP-адресам. Чтобы избежать блокировок вашей инфраструктуры, ознакомьтесь с практическими методами в статье: Как обойти блокировки по IP при парсинге стакана Raydium на Python.




