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 биржи. Это позволит вести детальный аудит всех операций без ущерба для производительности торговой системы.




