/
tsps
/
case_loader
Обзор
Документация
Войти
/
tsps
/
case_loader
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
master
app/services/job_processor.py
100 строк
4 KB
Василий Петров
Фикс бага - ошибка названия поля задачи
01 мар 2026, 20:55
01 мар 2026, 20:55
aad674d
Код
Авторство
О чём код?
import logging from datetime import datetime, timezone from sqlalchemy import select, update from sqlalchemy.engine import Connection from app.core.db import engine from app.services.pdf_parser import PdfParser from app.services.data_fetcher import DataFetcher from app.models.tables import document_import_job logger = logging.getLogger(__name__) def _get_job_status_and_params(conn: Connection, job_id: int) -> tuple[str, dict]: """ Читает статус и параметры поиска из таблицы document_import_job по ID задачи. Returns: Кортеж (статус, search_params) """ stmt = select( document_import_job.c.status, document_import_job.c.search_params ).where(document_import_job.c.id == job_id) result = conn.execute(stmt).fetchone() if not result: raise ValueError(f"Задача с ID {job_id} не найдена в базе данных") return result[0], result[1] # status, search_params def _update_job_status(conn: Connection, job_id: int, status: str, error_message: str | None = None, result_summary: dict | None = None): """ Обновляет статус задачи и, при необходимости, сообщение об ошибке. Args: conn: активное соединение с базой данных job_id: идентификатор задачи status: новый статус ('in_progress', 'completed', 'failed') error_message: текст ошибки (только для статуса 'failed') result_summary: итоговая статистика выполнения (только для статуса 'completed') """ update_values = { "status": status, "error_message": error_message # всегда перезаписываем, чтобы очистить старую ошибку при успехе } if status == "completed": update_values["updated_at"] = datetime.now(timezone.utc) if result_summary is not None: update_values["result_summary"] = result_summary stmt = ( update(document_import_job) .where(document_import_job.c.id == job_id) .values(**update_values) ) conn.execute(stmt) conn.commit() # Явный коммит, так как autocommit отключён def process_document_import_job(job_id: int) -> None: """ Фоновая задача RQ: обработка массовой загрузки документов. Args: job_id: идентификатор задачи в таблице document_import_job """ logger.info(f"Начата обработка задачи {job_id}") with engine.connect() as conn: try: # Читаем параметры поиска status, search_params = _get_job_status_and_params(conn, job_id) # Пропускаем задачи, которые уже обрабатываются или завершены if status not in ("created", "failed"): logger.info(f"Задача {job_id} имеет статус '{status}', пропускаем обработку") return logger.info(f"Получены параметры поиска для задачи {job_id}") # Обновляем статус на 'in_progress' _update_job_status(conn, job_id, "in_progress") logger.info(f"Статус задачи {job_id} изменён на 'in_progress'") # Выполняем задачу parser = PdfParser() fetcher = DataFetcher(parser=parser, db_conn=conn) result = fetcher.fetch_cases(search_params) # Завершаем задачу _update_job_status(conn, job_id, "completed", result_summary=result) logger.info(f"Задача {job_id} успешно завершена") except Exception as e: error_msg = f"Ошибка при выполнении задачи {job_id}: {str(e)}" logger.exception(error_msg) # Обновляем статус даже при ошибке — но вне try, чтобы избежать подавления исключения with engine.connect() as err_conn: _update_job_status(err_conn, job_id, "failed", error_message=str(e)) raise