/
tsps
/
ai_tools
Обзор
Документация
Войти
/
tsps
/
ai_tools
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
master
app/workers/summary_consumer.py
171 строка
8 KB
Василий Петров
Пропуск саммаризации для коротких текстов
04 май 2026, 21:34
04 май 2026, 21:34
52eb55e
Код
Авторство
О чём код?
import json import logging from redis import Redis from sqlalchemy import select, update, func from concurrent.futures import ThreadPoolExecutor from app.core.logging_config import setup_logging logger = logging.getLogger(__name__) logging.basicConfig( level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s" ) from app.core.config import config from app.core.db import get_db_connection from app.models.tables import audio_recording from app.services.summary import SummaryService MAX_RETRIES = 3 SUMM_MODEL_PATH = "ai_models/summary" # Синглтон Summary-сервиса _summary_service: SummaryService | None = None # Пул потоков для неблокирующего вызова модели (один на модуль) _summary_executor = ThreadPoolExecutor(max_workers=1) def _run_summary_sync(text: str) -> str: """ Обёртка для вызова SummaryService в потоке. Аргументы: text (str): Исходный текст для саммаризации. Возвращает: str: Сгенерированное краткое содержание. """ if _summary_service is None: raise RuntimeError("Summary-сервис не инициализирован") return _summary_service.summarize(text) def process_summarization_task(payload: dict) -> None: """ Читает transcription из audio_recording => Вызывает AI-инструменты для саммаризации => Сохраняет результат в transcript_summary. """ recording_id = payload.get("recording_id") if not recording_id: logger.error("Задача саммаризации без recording_id") return logger.info(f"Начало обработки: recording_id={recording_id}") try: with get_db_connection() as conn: # Читаем готовый текст транскрибации stmt = select(audio_recording.c.transcription).where( audio_recording.c.id == recording_id ) result = conn.execute(stmt).fetchone() if not result: logger.error(f"Запись не найдена: recording_id={recording_id}") return transcription_text = result[0] if not transcription_text: logger.warning(f"Пустая транскрибация для recording_id={recording_id}, пропускаю саммаризацию") return logger.info(f"recording_id={recording_id}, длина текста: {len(transcription_text)} символов") if len(transcription_text) < 500: summary_text = transcription_text logger.info(f"Саммаризация пропущена: {len(summary_text)} символов, recording_id={recording_id}") else: # Вызываем саммаризацию в отдельном потоке (не блокирует воркер) future = _summary_executor.submit(_run_summary_sync, transcription_text) summary_text = future.result() logger.info(f"Саммаризация получена: {len(summary_text)} символов, recording_id={recording_id}") # Обновляем результат: сохраняем саммари stmt = ( update(audio_recording) .where(audio_recording.c.id == recording_id) .values( transcript_summary=summary_text, updated_at=func.now(), ) ) conn.execute(stmt) conn.commit() logger.info(f"Саммаризация завершена: recording_id={recording_id}") except Exception as e: logger.exception(f"Ошибка при обработке recording_id={recording_id}: {e}") # Сохраняем ошибку в поле error_message try: with get_db_connection() as conn: stmt = ( update(audio_recording) .where(audio_recording.c.id == recording_id) .values( error_message=str(e), updated_at=func.now(), ) ) conn.execute(stmt) conn.commit() except Exception as rollback_error: logger.error( f"Не удалось сохранить ошибку для recording_id={recording_id}: {rollback_error}" ) def run_summary_consumer() -> None: """ Основной цикл воркера: блочное чтение из очереди 'summary'. При ошибке обработки — возврат задачи в очередь с инкрементом retry_count. """ global _summary_service # Инициализация Summary-сервиса (один раз при старте) _summary_service = SummaryService(SUMM_MODEL_PATH) _summary_service.load_model() logger.info(f"Модель саммаризации загружена из {SUMM_MODEL_PATH}") redis_conn = Redis.from_url(config.REDIS_URL) logger.info(f"Воркер саммаризации запущен, слушаю очередь 'summary' на {config.REDIS_URL}") while True: try: # blpop блокирует, пока не появится элемент или не истечёт таймаут result = redis_conn.blpop("summary", timeout=30) # blpop возвращает None при таймауте if result is None: continue # таймаут, продолжаем ждать _queue_name, raw_payload = result payload = json.loads(raw_payload) retry_count = payload.get("retry_count", 0) try: process_summarization_task(payload) except Exception as e: if retry_count >= MAX_RETRIES: logger.error( f"Задача recording_id={payload.get('recording_id')} превысила лимит попыток ({MAX_RETRIES}), отбрасываю. Ошибка: {e}" ) # TODO: здесь можно отправлять в dead-letter очередь else: payload["retry_count"] = retry_count + 1 redis_conn.rpush("summary", json.dumps(payload)) logger.warning( f"Ошибка обработки recording_id={payload.get('recording_id')}, попытка {retry_count + 1}/{MAX_RETRIES}. Задача возвращена в очередь. Ошибка: {e}" ) except json.JSONDecodeError as e: logger.error(f"Некорректный JSON в очереди: {raw_payload}. Ошибка: {e}") # Некорректный payload не возвращаем — он только засорит очередь except Exception as e: # Ошибки подключения к Redis и другие критические logger.exception(f"Критическая ошибка в цикле воркера: {e}") # Не выходим из цикла — даём шанс на восстановление соединения if __name__ == "__main__": setup_logging() run_summary_consumer()