Асинхронное логирование торговых операций в MongoDB на Python

Асинхронное логирование торговых операций в MongoDB на Python ИНВЕСТИЦИИ, ВКЛАДЫ и СБЕРЕЖЕНИЯ
Пошаговое руководство по настройке неблокирующего асинхронного логирования торговых операций бота в MongoDB на Python с использованием библиотеки Motor и asyncio.
Суть: Использование библиотеки motor совместно с asyncio.Queue позволяет реализовать неблокирующее логирование торговых операций в MongoDB. Это исключает задержки (latency) при отправке логов в БД, сохраняя высокую скорость исполнения ордеров торгового робота.

Исходный код

В алгоритмическом трейдинге критически важна каждая миллисекунда. Традиционное синхронное логирование в базу данных блокирует поток выполнения, что приводит к проскальзываниям (slippage) и задержкам при обработке рыночных данных.

В то время как для хранения быстрых котировок часто используется Интеграция бота-маркетмейкера на Python с Redis для обновления котировок, а для хранения сырых тиков — Интеграция TimescaleDB для хранения тиков с Binance на Python, структурированные логи торговых ордеров и транзакций идеально ложатся в документо-ориентированную базу данных MongoDB.

Ниже представлен готовый класс асинхронного логгера на базе паттерна Producer-Consumer, который гарантирует мгновенную запись в очередь и фоновую отправку данных в MongoDB.

import asyncio
import logging
from datetime import datetime
from typing import Dict, Any
from motor.motor_asyncio import AsyncIOMotorClient

class AsyncMongoTradingLogger:
    def __init__(self, mongo_uri: str, db_name: str, collection_name: str):
        self.client = AsyncIOMotorClient(mongo_uri)
        self.db = self.client[db_name]
        self.collection = self.db[collection_name]
        self.queue = asyncio.Queue()
        self._worker_task = None

    def start(self):
        """Запуск фонового воркера для обработки логов."""
        self._worker_task = asyncio.create_task(self._logger_worker())

    async def log_trade(self, order_id: str, symbol: str, side: str, price: float, qty: float, status: str, extra: Dict[str, Any] = None):
        """Быстрое добавление лога в очередь без блокировки основного потока."""
        log_entry = {
            "timestamp": datetime.utcnow(),
            "order_id": order_id,
            "symbol": symbol,
            "side": side,
            "price": price,
            "qty": qty,
            "status": status,
            "extra": extra or {}
        }
        await self.queue.put(log_entry)

    async def _logger_worker(self):
        """Фоновый воркер, записывающий логи в MongoDB."""
        while True:
            log_entry = await self.queue.get()
            try:
                await self.collection.insert_one(log_entry)
            except Exception as e:
                print(f"Ошибка записи лога в MongoDB: {e}")
            finally:
                self.queue.task_done()

    async def shutdown(self):
        """Корректное завершение работы с гарантией записи всех логов."""
        await self.queue.join()
        if self._worker_task:
            self._worker_task.cancel()
            try:
                await self._worker_task
            except asyncio.CancelledError:
                pass
        self.client.close()

# Пример использования в торговом цикле
async def main():
    logger = AsyncMongoTradingLogger("mongodb://localhost:27017", "trading_bot", "order_logs")
    logger.start()

    # Имитация отправки ордера
    print("Отправка ордера...")
    await logger.log_trade(
        order_id="order_102938",
        symbol="BTCUSDT",
        side="BUY",
        price=65250.5,
        qty=0.05,
        status="FILLED",
        extra={"strategy": "Grid_V1", "latency_ms": 12}
    )
    print("Ордер залогирован в очередь!")

    # Даем воркеру время обработать запись перед выходом
    await logger.shutdown()

if __name__ == "__main__":
    asyncio.run(main())

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

  • AsyncIOMotorClient: Асинхронный клиент из библиотеки motor, который обеспечивает неблокирующее сетевое взаимодействие с СУБД MongoDB в рамках общего event loop.
  • asyncio.Queue: Потокобезопасная асинхронная очередь. Метод log_trade мгновенно помещает лог в очередь, не дожидаясь ответа от базы данных, что сводит задержку логирования к нулю.
  • _logger_worker: Фоновая корутина, которая непрерывно извлекает элементы из очереди и записывает их в MongoDB.
  • shutdown(): Метод, гарантирующий корректное завершение работы. Он ожидает обработки всех оставшихся в очереди логов с помощью self.queue.join() перед закрытием соединения.

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

Для запуска данного решения вам потребуется установленная СУБД MongoDB и библиотека motor. Установите зависимость с помощью пакетного менеджера: pip install motor.

Запустите локальный экземпляр MongoDB или укажите строку подключения к вашему облачному кластеру в конструкторе класса AsyncMongoTradingLogger.

При интеграции логгера в реального торгового робота (например, при отправке подписанных запросов на биржу с использованием Генерации подписи HMAC SHA256 для Bybit API V5 на Python), вызывайте метод log_trade сразу после получения ответа от API биржи. Это позволит вести детальный аудит всех операций без ущерба для производительности торговой системы.

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