Интеграция InfluxDB и Grafana для HFT: Мониторинг Latency и Slippage

Интеграция InfluxDB и Grafana для HFT: Мониторинг Latency и Slippage ЛИЧНЫЙ БЮДЖЕТ и ЭКОНОМИЯ
Пошаговое руководство по интеграции InfluxDB и Grafana для мониторинга метрик HFT-бота в реальном времени. Настройка асинхронной записи Latency и Slippage на Python.
Суть: Для мониторинга HFT-ботов в реальном времени стандартные реляционные БД не подходят из-за высокой задержки на запись. Решением является использование InfluxDB 2.x с асинхронным пакетным клиентом на Python для сбора метрик задержки (Latency) и проскальзывания (Slippage), с последующей визуализацией в Grafana через язык запросов Flux.

Исходный код

При работе с высокочастотными стратегиями, где скорость обработки L2-стакана критична (подробнее в статье Ускорение обработки тиковых данных L2 (Orderbook): Polars против Pandas), задержка исполнения (Execution Latency) и проскальзывание (Slippage) являются ключевыми метриками здоровья торговой системы. В отличие от пакетной обработки исторических данных, где отлично подходит Apache Airflow для ETL криптосвечей: Оркестрация исторических данных, для real-time мониторинга HFT-инфраструктуры требуется специализированная TSDB (Time Series Database) с минимальным overhead на запись.

Для минимизации влияния логики логирования на основной поток исполнения HFT-бота, мы используем асинхронную пакетную запись (WriteOptions) в InfluxDB. Ниже представлен класс HFTMetricsCollector, оптимизированный для работы под высокой нагрузкой.

import time
import random
from influxdb_client import InfluxDBClient, Point, WriteOptions
from influxdb_client.client.write_api import ASYNCHRONOUS

class HFTMetricsCollector:
    def __init__(self, url: str, token: str, org: str, bucket: str):
        self.client = InfluxDBClient(url=url, token=token, org=org)
        # Настройка пакетной записи для минимизации I/O block
        self.write_api = self.client.write_api(write_options=WriteOptions(
            batch_size=500,           # Отправка пачками по 500 точек
            flush_interval=100,       # Сброс буфера каждые 100 мс
            jitter_interval=10,       # Случайная задержка для сглаживания пиков нагрузки
            retry_interval=1000,      # Интервал повтора при ошибке сети
            max_retries=3
        ))
        self.bucket = bucket
        self.org = org

    def log_execution(self, symbol: str, side: str, latency_ms: float, slippage_ticks: float, price: float, quantity: float):
        """
        Запись метрик исполнения ордера в InfluxDB.
        """
        point = Point("hft_execution") \
            .tag("symbol", symbol) \
            .tag("side", side) \
            .field("latency", float(latency_ms)) \
            .field("slippage", float(slippage_ticks)) \
            .field("price", float(price)) \
            .field("quantity", float(quantity)) \
            .time(time.time_ns())  # Наносекундная точность для HFT
        
        self.write_api.write(bucket=self.bucket, org=self.org, record=point)

    def close(self):
        self.write_api.close()
        self.client.close()

# Пример симуляции работы HFT-бота
if __name__ == "__main__":
    collector = HFTMetricsCollector(
        url="http://localhost:8086",
        token="my-super-secret-auth-token",
        org="quant_firm",
        bucket="hft_metrics"
    )

    print("Симуляция отправки метрик запущена...")
    try:
        while True:
            # Симулируем задержку сети + биржи (в миллисекундах)
            simulated_latency = random.uniform(0.8, 15.5)
            # Симулируем проскальзывание в тиках (может быть отрицательным)
            simulated_slippage = random.choice([-1.0, 0.0, 0.0, 1.0, 2.0, 5.0])
            
            collector.log_execution(
                symbol="BTCUSDT",
                side=random.choice(["BUY", "SELL"]),
                latency_ms=simulated_latency,
                slippage_ticks=simulated_slippage,
                price=65000.0 + random.uniform(-10, 10),
                quantity=0.05
            )
            time.sleep(0.01)  # Частота отправки - 100 ордеров в секунду
    except KeyboardInterrupt:
        collector.close()
        print("Сбор метрик остановлен.")

Сбор сырых метрик также требует последующей фильтрации аномалий сети, аналогично тому, как выполняется Очистка OHLCV данных в Pandas: Удаление выбросов и NaN, но уже на стороне Grafana с помощью Flux-запросов. Ниже представлен Flux-запрос для построения графика 99-го перцентиля задержки (Latency p99) в Grafana:

from(bucket: "hft_metrics")
  |> range(start: v.timeRangeStart, stop: v.timeRangeStop)
  |> filter(fn: (r) => r["_measurement"] == "hft_execution")
  |> filter(fn: (r) => r["_field"] == "latency")
  |> aggregateWindow(every: 1s, fn: (column, tables=<-) => tables |> quantile(q: 0.99), createEmpty: false)
  |> yield(name: "p99_latency")

Разбор параметров

  • batch_size: Размер пачки данных для отправки. Значение 500 оптимально для снижения сетевого оверхеда без создания критической задержки отображения.
  • flush_interval: Максимальное время ожидания (в мс) перед принудительной отправкой неполного батча. Значение 100 гарантирует real-time отображение в Grafana.
  • time_ns(): Использование наносекундного таймстампа критически важно для HFT, так как в одну миллисекунду может уложиться несколько транзакций.
  • quantile(q: 0.99): Позволяет отслеживать «хвосты» распределения задержки (tail latency), что критично для оценки стабильности подключения к бирже.

Как запустить

1. Разверните стек InfluxDB и Grafana с помощью Docker Compose:

version: '3.8'
services:
  influxdb:
    image: influxdb:2.7
    ports:
      - "8086:8086"
    environment:
      - DOCKER_INFLUXDB_INIT_MODE=setup
      - DOCKER_INFLUXDB_INIT_USERNAME=admin
      - DOCKER_INFLUXDB_INIT_PASSWORD=my-super-secret-password
      - DOCKER_INFLUXDB_INIT_ORG=quant_firm
      - DOCKER_INFLUXDB_INIT_BUCKET=hft_metrics
      - DOCKER_INFLUXDB_INIT_ADMIN_TOKEN=my-super-secret-auth-token
  grafana:
    image: grafana/grafana:10.0.0
    ports:
      - "3000:3000"
    depends_on:
      - influxdb

2. Запустите контейнеры командой docker-compose up -d.

3. Установите библиотеку клиента InfluxDB для Python: pip install influxdb-client.

4. Запустите Python-скрипт для генерации тестовой нагрузки.

5. Откройте Grafana (http://localhost:3000, логин/пароль по умолчанию: admin/admin), добавьте источник данных (Data Source) типа InfluxDB, выберите язык запросов Flux, укажите URL http://influxdb:8086, организацию quant_firm и токен из конфигурационного файла.

6. Создайте новый Dashboard, добавьте панель Time Series и вставьте Flux-запрос для визуализации Latency или Slippage.

Оцените статью
FinFluct