/
Anna_brv
/
ASR-Service
Обзор
Документация
Войти
/
Anna_brv
/
ASR-Service
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
src/worker.py
89 строк
3 KB
Anna
first_commit
30 апр 2026, 21:52
30 апр 2026, 21:52
283885c
Код
Авторство
О чём код?
#!/usr/bin/env python3 """RQ Worker script for ASR service.""" import signal import socket import uuid import sys import os import redis from rq import Worker from src.config import settings from loguru import logger # Using the same logger as the existing code import multiprocessing as mp mp.set_start_method("spawn", force=True) WORKER_NAME_PREFIX = "asr-worker" class GracefulKiller: """Handle SIGTERM and SIGINT for graceful shutdown.""" kill_now = False def __init__(self): signal.signal(signal.SIGTERM, self._handle_exit) signal.signal(signal.SIGINT, self._handle_exit) def _handle_exit(self, signum, frame): logger.warning(f"Received signal {signum}, shutting down gracefully...") self.kill_now = True def cleanup_stale_workers(redis_conn: redis.Redis, prefix: str) -> None: """Clean up stale workers that are no longer running.""" all_workers = Worker.all(connection=redis_conn) for w in all_workers: if w.name.startswith(prefix): try: # Simple check: if worker state is not 'idle' or 'busy', it might be stale # A more robust check would involve verifying the PID on the host if w.state not in ("idle", "busy"): logger.info(f"Cleaning up stale worker: {w.name} (state: {w.state})") w.register_death() except Exception as e: logger.exception(f"Error checking worker {w.name}: {e}") def main(): """Run RQ worker.""" GracefulKiller() try: redis_url = settings.redis_url logger.info(f"Connecting to Redis at {redis_url}") redis_conn = redis.from_url(redis_url) redis_conn.ping() # Verify connection logger.info("Connected to Redis successfully") # Clean up potentially stale workers before starting cleanup_stale_workers(redis_conn, WORKER_NAME_PREFIX) # Generate a unique worker name unique_suffix = uuid.uuid4().hex[:8] hostname = socket.gethostname() worker_name = f"{WORKER_NAME_PREFIX}-{hostname}-{unique_suffix}" logger.info(f"Using worker name: {worker_name}") # Use configurable queues queues = settings.worker_queues logger.info(f"Listening on queues: {queues}") worker = Worker(queues, connection=redis_conn, name=worker_name) logger.info("Starting RQ worker for ASR service...") worker.work(logging_level=settings.log_level.upper()) except redis.ConnectionError: logger.error("Could not connect to Redis") sys.exit(1) except KeyboardInterrupt: logger.info("Worker interrupted by user") except Exception as e: logger.exception(f"Failed to start RQ worker: {e}") sys.exit(1) if __name__ == "__main__": main()