x
This commit is contained in:
+16
-2
@@ -11,8 +11,14 @@ from utils.file_handler import FileHandler
|
|||||||
from services.storage_integration_rustfs import StorageIntegrationService
|
from services.storage_integration_rustfs import StorageIntegrationService
|
||||||
from services.redis_task_manager import redis_task_manager
|
from services.redis_task_manager import redis_task_manager
|
||||||
from services.task_query_service import TaskQueryService
|
from services.task_query_service import TaskQueryService
|
||||||
from celery_tasks import process_stp_task
|
|
||||||
from database.database import get_db_session
|
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 utils.logger import get_logger
|
||||||
from sqlalchemy.ext.asyncio import AsyncSession
|
from sqlalchemy.ext.asyncio import AsyncSession
|
||||||
from core.cavity_layout_optimizer import CavityLayoutOptimizer
|
from core.cavity_layout_optimizer import CavityLayoutOptimizer
|
||||||
@@ -103,8 +109,16 @@ async def upload_stp(
|
|||||||
|
|
||||||
# 后台处理走统一编排服务,避免请求会话在后台失效
|
# 后台处理走统一编排服务,避免请求会话在后台失效
|
||||||
process_params = {"material": material}
|
process_params = {"material": material}
|
||||||
|
if _use_celery:
|
||||||
process_stp_task.delay(task_id, str(file_path), stp_file.id, process_params)
|
process_stp_task.delay(task_id, str(file_path), stp_file.id, process_params)
|
||||||
logger.info(f"[UPLOAD] 后台处理已调度: task_id={task_id}")
|
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 {
|
return {
|
||||||
"task_id": task_id,
|
"task_id": task_id,
|
||||||
|
|||||||
@@ -13,7 +13,13 @@ from utils.logger import get_logger
|
|||||||
from sqlalchemy.ext.asyncio import AsyncSession
|
from sqlalchemy.ext.asyncio import AsyncSession
|
||||||
from services.auth_service import get_current_active_user
|
from services.auth_service import get_current_active_user
|
||||||
from models.database import User
|
from models.database import User
|
||||||
|
|
||||||
|
try:
|
||||||
from celery_tasks import process_stp_task
|
from celery_tasks import process_stp_task
|
||||||
|
_use_celery = True
|
||||||
|
except ImportError:
|
||||||
|
process_stp_task = None
|
||||||
|
_use_celery = False
|
||||||
|
|
||||||
logger = get_logger(__name__)
|
logger = get_logger(__name__)
|
||||||
|
|
||||||
@@ -90,8 +96,16 @@ async def upload_stp(
|
|||||||
task_info["file_hash"] = file_meta["sha256"]
|
task_info["file_hash"] = file_meta["sha256"]
|
||||||
await redis_task_manager.set_task(task_id, task_info)
|
await redis_task_manager.set_task(task_id, task_info)
|
||||||
|
|
||||||
|
if _use_celery:
|
||||||
process_stp_task.delay(task_id, str(file_path), stp_file.id, process_params)
|
process_stp_task.delay(task_id, str(file_path), stp_file.id, process_params)
|
||||||
logger.info(f"[UPLOAD] 后台处理已调度: task_id={task_id}")
|
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 {
|
return {
|
||||||
"task_id": task_id,
|
"task_id": task_id,
|
||||||
|
|||||||
Reference in New Issue
Block a user