/
tegdif
/
cism
Обзор
Документация
Войти
/
tegdif
/
cism
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
src/workers/consumer.py
144 строки
5 KB
tegdif
first_commit
28 июл 2026, 08:54
28 июл 2026, 08:54
8ecbd07
Код
Авторство
О чём код?
import asyncio import json import logging import signal import uuid from datetime import datetime import aio_pika from src.core.config import settings from src.core.database import async_session_factory from src.core.rabbitmq import rabbitmq_client from src.models.task import TaskStatus, TaskPriority from src.services.task_service import TaskService logger = logging.getLogger(__name__) # Порядок приоритетов: сначала обрабатываются HIGH задачи PRIORITY_ORDER = [TaskPriority.HIGH, TaskPriority.MEDIUM, TaskPriority.LOW] async def process_task(task_id: uuid.UUID) -> None: """Обрабатывает одну задачу: переход PENDING -> IN_PROGRESS -> COMPLETED/FAILED.""" async with async_session_factory() as session: service = TaskService(session) task = await service.get_task(task_id) if not task: logger.error("Задача %s не найдена", task_id) return if task.status != TaskStatus.PENDING: logger.info("Задача %s не в статусе PENDING (текущий статус=%s), пропускаем", task_id, task.status.value) return await service.update_task_status(task_id, TaskStatus.IN_PROGRESS) logger.info("Задача %s начала обработку", task_id) try: # Имитация реальной работы await asyncio.sleep(0.5) result = f"Задача '{task.title}' успешно обработана в {datetime.utcnow().isoformat()}" await service.update_task_status( task_id, TaskStatus.COMPLETED, result=result ) logger.info("Задача %s успешно завершена", task_id) except Exception as e: error_info = f"Ошибка обработки задачи {task_id}: {str(e)}" await service.update_task_status( task_id, TaskStatus.FAILED, error_info=error_info ) logger.error(error_info) async def on_message(message: aio_pika.IncomingMessage) -> None: """Обрабатывает входящее сообщение из любой очереди приоритетов.""" async with message.process(ignore_processed=True): try: body = json.loads(message.body.decode()) task_id = uuid.UUID(body["task_id"]) async with async_session_factory() as session: service = TaskService(session) task = await service.get_task(task_id) if task and task.status == TaskStatus.NEW: await service.update_task_status(task_id, TaskStatus.PENDING) await process_task(task_id) except Exception as e: logger.error("Ошибка обработки сообщения: %s", e, exc_info=True) async def start_worker() -> None: """Запускает worker, читающий сообщения из всех очередей приоритетов.""" logger.info("Запуск worker'а задач...") await rabbitmq_client.connect() semaphore = asyncio.Semaphore(settings.max_concurrent_tasks) # Флаг для graceful shutdown shutdown_event = asyncio.Event() def _signal_handler() -> None: logger.info("Получен сигнал остановки, завершаем worker...") shutdown_event.set() loop = asyncio.get_event_loop() for sig in (signal.SIGINT, signal.SIGTERM): loop.add_signal_handler(sig, _signal_handler) async def _process_with_semaphore(msg: aio_pika.IncomingMessage) -> None: async with semaphore: await on_message(msg) async def _consume_queue(priority: TaskPriority) -> None: """Читает сообщения из одной очереди приоритетов.""" queue = await rabbitmq_client.consume(priority) logger.info( "Чтение из очереди '%s' (приоритет=%s)", settings.queue_names[priority.value], priority.value, ) async with queue.iterator() as queue_iter: async for message in queue_iter: if shutdown_event.is_set(): break asyncio.create_task(_process_with_semaphore(message)) # Запускаем consumer'ы для каждой очереди приоритетов по порядку tasks = [] for priority in PRIORITY_ORDER: task = asyncio.create_task(_consume_queue(priority)) tasks.append(task) logger.info( "Worker запущен с %d очередями приоритетов, max_concurrent=%d", len(PRIORITY_ORDER), settings.max_concurrent_tasks, ) # Ожидаем сигнала остановки await shutdown_event.wait() logger.info("Worker завершает работу...") # Отменяем задачи consumer'ов for t in tasks: t.cancel() await asyncio.gather(*tasks, return_exceptions=True) await rabbitmq_client.close() logger.info("Worker полностью остановлен.") if __name__ == "__main__": logging.basicConfig( level=logging.INFO, format="%(asctime)s - %(name)s - %(levelname)s - %(message)s", ) try: asyncio.run(start_worker()) except KeyboardInterrupt: logger.info("Worker остановлен пользователем")