/
Arslan16
/
FleetHub
Обзор
Документация
Войти
/
Arslan16
/
FleetHub
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
dev
services/monitoring/utils.py
80 строк
3 KB
Arslan16
Заработала запись метрик в TimescaleDB из rmq. Проект переведен с PostgreSQL 15 на PostgreSQL 17
11 авг 2026, 19:42
11 авг 2026, 19:42
83f4eae
Код
Авторство
О чём код?
import asyncio import aio_pika from aiormq import AMQPConnectionError from config import RABBITMQ_ASYNC_SESSIONMAKER from database_utils import MetricsManager from logger import logger from rmq_manager import RABBITMQ_MANAGER, RABBITMQ_METRICS_QUEUE from schemas import RabbitMQQueueData async def cancel_async_task(task: asyncio.Task | None): """ Отменяет переданную асинхронную задачу и ожидает её завершения.\n Штатная ошибка `asyncio.CancelledError`, возникающая в результате\n отмены задачи, перехватывается и выводится в стандартный вывод.\n Args: task (asyncio.Task | None): Асинхронная задача, которую необходимо отменить. Если передано значение `None` или объект, не являющийся `asyncio.Task`, функция ничего не выполняет. """ if isinstance(task, asyncio.Task): task.cancel() try: await task except asyncio.CancelledError as exc: print(exc) async def accept_metrics( message: aio_pika.IncomingMessage ): try: async with message.process(): async with RABBITMQ_ASYNC_SESSIONMAKER() as session: manager = MetricsManager(session) rmq_data: RabbitMQQueueData = RabbitMQQueueData.model_validate_json(message.body) for code, value in rmq_data.metrics.model_dump().items(): await manager.create( agent_id=rmq_data.agent_id, code=code, value=value ) await session.commit() except asyncio.CancelledError: # Фоновые таски могут быть отменены при shutdown loop logger.warning("Callback on_message был отменён") raise except Exception as exc: # Любые другие ошибки при обработке сообщения logger.exception(exc) raise async def run_try_rabbitmq_connection_loop(): """ Ожидает доступность RabbitMQ и выполняет повторные попытки подключения. При недоступности RabbitMQ функция циклически вызывает метод подключения менеджера RabbitMQ с интервалом в 2 секунды до тех пор, пока соединение не будет установлено. Raises: AMQPConnectionError: Перехватывается внутри функции. Ошибка логируется, после чего выполняется повторная попытка подключения. """ while not RABBITMQ_MANAGER.is_available(): try: await RABBITMQ_MANAGER.connect() except AMQPConnectionError as exc: logger.exception(exc) await asyncio.sleep(2) async def register_callback_on_metrics_queue(): await run_try_rabbitmq_connection_loop() await RABBITMQ_MANAGER.register_callback_on_queue(RABBITMQ_METRICS_QUEUE, accept_metrics)