Исходный код
При работе с высокочастотными стратегиями, где скорость обработки 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.




