From 3443763fbf7de5dc90fa01282e8a3e436d05e63a Mon Sep 17 00:00:00 2001 From: chenjw28 <792430652@qq.com> Date: Fri, 29 May 2026 09:47:21 +0800 Subject: [PATCH] x --- src/api/routes.py | 20 +++++++++++++++++--- src/api/v1/upload_router.py | 20 +++++++++++++++++--- 2 files changed, 34 insertions(+), 6 deletions(-) diff --git a/src/api/routes.py b/src/api/routes.py index f5af82f..4810fcf 100644 --- a/src/api/routes.py +++ b/src/api/routes.py @@ -11,8 +11,14 @@ from utils.file_handler import FileHandler from services.storage_integration_rustfs import StorageIntegrationService from services.redis_task_manager import redis_task_manager from services.task_query_service import TaskQueryService -from celery_tasks import process_stp_task from database.database import get_db_session + +try: + from celery_tasks import process_stp_task + _use_celery = True +except ImportError: + process_stp_task = None + _use_celery = False from utils.logger import get_logger from sqlalchemy.ext.asyncio import AsyncSession from core.cavity_layout_optimizer import CavityLayoutOptimizer @@ -103,8 +109,16 @@ async def upload_stp( # 后台处理走统一编排服务,避免请求会话在后台失效 process_params = {"material": material} - process_stp_task.delay(task_id, str(file_path), stp_file.id, process_params) - logger.info(f"[UPLOAD] 后台处理已调度: task_id={task_id}") + if _use_celery: + process_stp_task.delay(task_id, str(file_path), stp_file.id, process_params) + logger.info(f"[UPLOAD] Celery 任务已调度: task_id={task_id}") + else: + import asyncio + from services.processing_service import processing_service + asyncio.create_task(processing_service.process_file_with_storage( + task_id, str(file_path), stp_file.id, process_params + )) + logger.info(f"[UPLOAD] 直接后台处理: task_id={task_id} (celery 未安装)") return { "task_id": task_id, diff --git a/src/api/v1/upload_router.py b/src/api/v1/upload_router.py index 1b6888b..8999a39 100644 --- a/src/api/v1/upload_router.py +++ b/src/api/v1/upload_router.py @@ -13,7 +13,13 @@ from utils.logger import get_logger from sqlalchemy.ext.asyncio import AsyncSession from services.auth_service import get_current_active_user from models.database import User -from celery_tasks import process_stp_task + +try: + from celery_tasks import process_stp_task + _use_celery = True +except ImportError: + process_stp_task = None + _use_celery = False logger = get_logger(__name__) @@ -90,8 +96,16 @@ async def upload_stp( task_info["file_hash"] = file_meta["sha256"] await redis_task_manager.set_task(task_id, task_info) - process_stp_task.delay(task_id, str(file_path), stp_file.id, process_params) - logger.info(f"[UPLOAD] 后台处理已调度: task_id={task_id}") + if _use_celery: + process_stp_task.delay(task_id, str(file_path), stp_file.id, process_params) + logger.info(f"[UPLOAD] Celery 任务已调度: task_id={task_id}") + else: + import asyncio + from services.processing_service import processing_service + asyncio.create_task(processing_service.process_file_with_storage( + task_id, str(file_path), stp_file.id, process_params + )) + logger.info(f"[UPLOAD] 直接后台处理: task_id={task_id} (celery 未安装)") return { "task_id": task_id,